fix(playback): prune downloaded transcode segments
This commit is contained in:
@@ -243,10 +243,11 @@ func NewPlaybackHandler(sessionMgr SessionManagerInterface, opts ...FilePathReso
|
||||
h.tm.Config = func() playback.TranscodeRuntimeConfig {
|
||||
c := h.playbackConfig()
|
||||
return playback.TranscodeRuntimeConfig{
|
||||
TranscodeDir: c.TranscodeDir,
|
||||
FFmpegPath: c.FFmpegPath,
|
||||
HWAccel: c.HWAccel,
|
||||
HWDevice: c.HWDevice,
|
||||
TranscodeDir: c.TranscodeDir,
|
||||
FFmpegPath: c.FFmpegPath,
|
||||
HWAccel: c.HWAccel,
|
||||
HWDevice: c.HWDevice,
|
||||
SegmentRetentionSeconds: c.SegmentRetentionSeconds,
|
||||
}
|
||||
}
|
||||
h.tm.StartThrottler = func(ctx context.Context, ts *playback.TranscodeSession) {
|
||||
@@ -322,8 +323,9 @@ func (h *PlaybackHandler) playbackConfig() config.PlaybackConfig {
|
||||
return h.PlaybackConfig()
|
||||
}
|
||||
return config.PlaybackConfig{
|
||||
TranscodeEnabled: true,
|
||||
TranscodeDir: filepath.Join(os.TempDir(), "silo-transcode"),
|
||||
TranscodeEnabled: true,
|
||||
TranscodeDir: filepath.Join(os.TempDir(), "silo-transcode"),
|
||||
SegmentRetentionSeconds: 600,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3446,32 +3448,33 @@ func (h *PlaybackHandler) HandleStartTranscode(w http.ResponseWriter, r *http.Re
|
||||
}
|
||||
|
||||
localOpts := playback.TranscodeOpts{
|
||||
InputPath: file.FilePath,
|
||||
OutputDir: filepath.Join(playbackCfg.TranscodeDir, req.SessionID),
|
||||
SessionID: req.SessionID,
|
||||
SourceVideoCodec: file.CodecVideo,
|
||||
VideoBitstreamFilter: videoBitstreamFilter,
|
||||
SeekSeconds: transportSeekSeconds,
|
||||
StreamOriginSeconds: streamOriginSeconds,
|
||||
CopySeekAnchorResolved: videoCopy,
|
||||
StartSegmentNumber: startSegmentNumber,
|
||||
TargetResolution: req.TargetResolution,
|
||||
TargetCodecVideo: req.TargetCodecVideo,
|
||||
TargetCodecAudio: req.TargetCodecAudio,
|
||||
TargetBitrateKbps: req.TargetBitrateKbps,
|
||||
SegmentDuration: req.SegmentDuration,
|
||||
FFmpegPath: playbackCfg.FFmpegPath,
|
||||
HWAccel: playbackCfg.HWAccel,
|
||||
HWDevice: playbackCfg.HWDevice,
|
||||
AudioTrackIndex: session.AudioTrackIndex,
|
||||
SubtitleTrackIndex: req.SubtitleTrackIndex,
|
||||
SubtitleBurnIn: req.SubtitleBurnIn,
|
||||
SubtitleCodec: subtitleCodec,
|
||||
TotalDuration: float64(file.Duration),
|
||||
FastStart: true,
|
||||
NodeType: "integrated",
|
||||
ExecutionMode: "integrated",
|
||||
FFmpegLogSink: h.FFmpegLogSink,
|
||||
InputPath: file.FilePath,
|
||||
OutputDir: filepath.Join(playbackCfg.TranscodeDir, req.SessionID),
|
||||
SessionID: req.SessionID,
|
||||
SourceVideoCodec: file.CodecVideo,
|
||||
VideoBitstreamFilter: videoBitstreamFilter,
|
||||
SeekSeconds: transportSeekSeconds,
|
||||
StreamOriginSeconds: streamOriginSeconds,
|
||||
CopySeekAnchorResolved: videoCopy,
|
||||
StartSegmentNumber: startSegmentNumber,
|
||||
TargetResolution: req.TargetResolution,
|
||||
TargetCodecVideo: req.TargetCodecVideo,
|
||||
TargetCodecAudio: req.TargetCodecAudio,
|
||||
TargetBitrateKbps: req.TargetBitrateKbps,
|
||||
SegmentDuration: req.SegmentDuration,
|
||||
SegmentRetentionSeconds: playbackCfg.SegmentRetentionSeconds,
|
||||
FFmpegPath: playbackCfg.FFmpegPath,
|
||||
HWAccel: playbackCfg.HWAccel,
|
||||
HWDevice: playbackCfg.HWDevice,
|
||||
AudioTrackIndex: session.AudioTrackIndex,
|
||||
SubtitleTrackIndex: req.SubtitleTrackIndex,
|
||||
SubtitleBurnIn: req.SubtitleBurnIn,
|
||||
SubtitleCodec: subtitleCodec,
|
||||
TotalDuration: float64(file.Duration),
|
||||
FastStart: true,
|
||||
NodeType: "integrated",
|
||||
ExecutionMode: "integrated",
|
||||
FFmpegLogSink: h.FFmpegLogSink,
|
||||
}
|
||||
startState := transcodeStartState{
|
||||
req: req,
|
||||
@@ -3680,7 +3683,8 @@ func (h *PlaybackHandler) HandleGetTranscodeSegment(w http.ResponseWriter, r *ht
|
||||
h.touchSessionActivity(sessionID)
|
||||
|
||||
segmentName := chi.URLParam(r, "name")
|
||||
segmentPath, err := transcodeSession.GetSegment(segmentName)
|
||||
downloadGeneration := transcodeSession.SegmentGeneration()
|
||||
segmentFile, segmentInfo, err := transcodeSession.OpenSegment(segmentName)
|
||||
if err != nil && errors.Is(err, playback.ErrSegmentNotFound) {
|
||||
segNum, parseErr := playback.ParseSegmentNumber(segmentName)
|
||||
if parseErr == nil {
|
||||
@@ -3717,7 +3721,7 @@ func (h *PlaybackHandler) HandleGetTranscodeSegment(w http.ResponseWriter, r *ht
|
||||
"session", sessionID,
|
||||
"playback_session_id", sessionID,
|
||||
)
|
||||
segmentPath, err = transcodeSession.WaitForSegment(segmentName, decision.WaitTimeout)
|
||||
segmentFile, segmentInfo, err = transcodeSession.WaitForOpenSegment(segmentName, decision.WaitTimeout)
|
||||
if err != nil && errors.Is(err, playback.ErrSegmentNotFound) {
|
||||
slog.InfoContext(r.Context(), "transcode segment wait timeout", "component", "api",
|
||||
"segment", segmentName,
|
||||
@@ -3779,7 +3783,7 @@ func (h *PlaybackHandler) HandleGetTranscodeSegment(w http.ResponseWriter, r *ht
|
||||
); restartErr == nil {
|
||||
// Throttler + exit monitor re-arm via the session's
|
||||
// restart hook.
|
||||
segmentPath, err = transcodeSession.WaitForSegment(segmentName, 30*time.Second)
|
||||
segmentFile, segmentInfo, err = transcodeSession.WaitForOpenSegment(segmentName, 30*time.Second)
|
||||
if err == nil && strings.EqualFold(transcodeSession.Opts().TargetCodecVideo, "copy") {
|
||||
// Copy-mode seeks can resume as soon as the target segment
|
||||
// exists, but that sometimes leaves the player one segment
|
||||
@@ -3787,7 +3791,9 @@ func (h *PlaybackHandler) HandleGetTranscodeSegment(w http.ResponseWriter, r *ht
|
||||
// for a single lookahead fragment when available so the
|
||||
// first resumed playback window is less brittle.
|
||||
nextSegmentName := fmt.Sprintf("seg_%05d.m4s", segNum+1)
|
||||
_, _ = transcodeSession.WaitForSegment(nextSegmentName, 1200*time.Millisecond)
|
||||
if nextSegment, _, nextErr := transcodeSession.WaitForOpenSegment(nextSegmentName, 1200*time.Millisecond); nextErr == nil {
|
||||
_ = nextSegment.Close()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -3795,7 +3801,7 @@ func (h *PlaybackHandler) HandleGetTranscodeSegment(w http.ResponseWriter, r *ht
|
||||
} else if transcodeSession.IsRunning() {
|
||||
// Non-numbered segment (e.g., init.mp4 for fMP4 HLS).
|
||||
// Wait briefly — the init segment is written almost immediately.
|
||||
segmentPath, err = transcodeSession.WaitForSegment(segmentName, 10*time.Second)
|
||||
segmentFile, segmentInfo, err = transcodeSession.WaitForOpenSegment(segmentName, 10*time.Second)
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
@@ -3807,14 +3813,19 @@ func (h *PlaybackHandler) HandleGetTranscodeSegment(w http.ResponseWriter, r *ht
|
||||
return
|
||||
}
|
||||
|
||||
// Report segment download for throttle tracking.
|
||||
if segNum, parseErr := playback.ParseSegmentNumber(segmentName); parseErr == nil {
|
||||
transcodeSession.ReportSegmentDownloaded(segNum)
|
||||
}
|
||||
|
||||
w.Header().Set("Cache-Control", "no-store, max-age=0")
|
||||
w.Header().Set("Pragma", "no-cache")
|
||||
http.ServeFile(w, r, segmentPath)
|
||||
defer func() { _ = segmentFile.Close() }()
|
||||
sw := httpstream.NewRollingDeadlineWriter(w)
|
||||
http.ServeContent(sw, r, segmentInfo.Name(), segmentInfo.ModTime(), segmentFile)
|
||||
if r.Method == http.MethodGet &&
|
||||
sw.StatusCode() == http.StatusOK &&
|
||||
sw.BytesWritten() == segmentInfo.Size() &&
|
||||
sw.Outcome(r.Context()) == httpstream.OutcomeCompleted {
|
||||
if segNum, parseErr := playback.ParseSegmentNumber(segmentName); parseErr == nil {
|
||||
transcodeSession.ReportSegmentDownloadedForGeneration(segNum, downloadGeneration)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// buildProxyManifestURL signs a stream token carrying the session's full
|
||||
|
||||
@@ -736,7 +736,7 @@ func (h *PlaybackHandler) prepareLocalTransportV3(r *http.Request, session *play
|
||||
sourceVideoCodec, sourceDuration := sourceExecutionMetadataV3(file, result)
|
||||
seekSeconds, startSegment := configureHLSTimelineV3(result.Plan, videoCodec, 2, sourceDuration)
|
||||
unlock := h.tm.LockSessionLifecycle(session.ID)
|
||||
ts, err := h.startLocalPlaybackTransport(r.Context(), playback.TranscodeOpts{InputPath: file.FilePath, OutputDir: outputDir, OutputSubdir: outputSubdir, SessionID: session.ID, SourceVideoCodec: sourceVideoCodec, VideoBitstreamFilter: videoBitstreamFilterForPlanV3(result.Plan), SeekSeconds: seekSeconds, StartSegmentNumber: startSegment, TargetResolution: result.TargetResolution, TargetCodecVideo: videoCodec, TargetCodecAudio: result.TargetAudioCodec, TargetAudioChannels: result.TargetAudioChannels, TargetBitrateKbps: result.TargetBitrateKbps, SegmentDuration: 2, FFmpegPath: cfg.FFmpegPath, HWAccel: cfg.HWAccel, HWDevice: cfg.HWDevice, AudioTrackIndex: plannedAudioTrackIndexV3(result, session.AudioTrackIndex), SubtitleTrackIndex: result.SubtitleTransportTrackIndex, SubtitleBurnIn: result.SubtitleBurnIn, SubtitleCodec: result.SubtitleCodec, TotalDuration: sourceDuration, FastStart: true, NodeType: playbackNodeIntegratedV3, ExecutionMode: playbackNodeIntegratedV3, FFmpegLogSink: h.FFmpegLogSink})
|
||||
ts, err := h.startLocalPlaybackTransport(r.Context(), playback.TranscodeOpts{InputPath: file.FilePath, OutputDir: outputDir, OutputSubdir: outputSubdir, SessionID: session.ID, SourceVideoCodec: sourceVideoCodec, VideoBitstreamFilter: videoBitstreamFilterForPlanV3(result.Plan), SeekSeconds: seekSeconds, StartSegmentNumber: startSegment, TargetResolution: result.TargetResolution, TargetCodecVideo: videoCodec, TargetCodecAudio: result.TargetAudioCodec, TargetAudioChannels: result.TargetAudioChannels, TargetBitrateKbps: result.TargetBitrateKbps, SegmentDuration: 2, SegmentRetentionSeconds: cfg.SegmentRetentionSeconds, FFmpegPath: cfg.FFmpegPath, HWAccel: cfg.HWAccel, HWDevice: cfg.HWDevice, AudioTrackIndex: plannedAudioTrackIndexV3(result, session.AudioTrackIndex), SubtitleTrackIndex: result.SubtitleTransportTrackIndex, SubtitleBurnIn: result.SubtitleBurnIn, SubtitleCodec: result.SubtitleCodec, TotalDuration: sourceDuration, FastStart: true, NodeType: playbackNodeIntegratedV3, ExecutionMode: playbackNodeIntegratedV3, FFmpegLogSink: h.FFmpegLogSink})
|
||||
if err != nil {
|
||||
unlock()
|
||||
return preparedTransportV3{}, &transportErrorV3{reason: "transcode_start_failed", message: "Failed to start the playback transport.", retryable: true, cause: err}
|
||||
|
||||
@@ -13,7 +13,10 @@ import (
|
||||
"github.com/robfig/cron/v3"
|
||||
)
|
||||
|
||||
const cloudflareURLMode = "cloudflare_token"
|
||||
const (
|
||||
cloudflareURLMode = "cloudflare_token"
|
||||
playbackSegmentRetentionSettingKey = "playback.segment_retention_seconds"
|
||||
)
|
||||
|
||||
// adminSettingDefaults is the effective value shown by the Admin UI when no
|
||||
// row exists in server_settings. Keep these values aligned with the runtime
|
||||
@@ -50,6 +53,7 @@ var adminSettingDefaults = map[string]string{
|
||||
|
||||
"playback.ffmpeg_path": "/usr/lib/jellyfin-ffmpeg/ffmpeg",
|
||||
"playback.transcode_dir": DefaultTranscodeDir,
|
||||
playbackSegmentRetentionSettingKey: "600",
|
||||
"playback.hw_accel": "auto",
|
||||
"playback.transcode_enabled": "true",
|
||||
"playback.local_transcode_fallback": "true",
|
||||
@@ -298,6 +302,16 @@ func NormalizeAdminSetting(key, raw string) (string, error) {
|
||||
return normalizeAdminInt(key, value, 1, 99)
|
||||
case "transcode_throttle_seconds":
|
||||
return normalizeAdminInt(key, value, 60, 86400)
|
||||
case playbackSegmentRetentionSettingKey:
|
||||
normalized, err := normalizeAdminInt(key, value, 0, 86400)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
seconds, _ := strconv.Atoi(normalized)
|
||||
if seconds != 0 && seconds < 120 {
|
||||
return "", fmt.Errorf("%s must be 0 or between 120 and 86400", key)
|
||||
}
|
||||
return normalized, nil
|
||||
case "ai.max_concurrent_jobs", "subtitle_ai.max_concurrent_jobs":
|
||||
return normalizeAdminInt(key, value, 1, 1024)
|
||||
case "subtitle_ai.batch_size":
|
||||
|
||||
@@ -217,6 +217,7 @@ func TestNormalizeAdminSettingRejectsInvalidValues(t *testing.T) {
|
||||
{key: "theme.catalog_url", value: "http://raw.githubusercontent.com/Silo-Server/silo-themes/main/catalog.json"},
|
||||
{key: "theme.catalog_url", value: "https://example.com/catalog.json"},
|
||||
{key: "redis.url", value: "not-a-url"},
|
||||
{key: "playback.segment_retention_seconds", value: "119"},
|
||||
}
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.key, func(t *testing.T) {
|
||||
@@ -227,6 +228,20 @@ func TestNormalizeAdminSettingRejectsInvalidValues(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeAdminSettingAcceptsSegmentRetentionBounds(t *testing.T) {
|
||||
for _, value := range []string{"0", "120", "86400"} {
|
||||
t.Run(value, func(t *testing.T) {
|
||||
got, err := NormalizeAdminSetting("playback.segment_retention_seconds", value)
|
||||
if err != nil {
|
||||
t.Fatalf("NormalizeAdminSetting: %v", err)
|
||||
}
|
||||
if got != value {
|
||||
t.Fatalf("normalized retention = %q, want %q", got, value)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeAdminSettingAcceptsApprovedThemeCatalogURL(t *testing.T) {
|
||||
got, err := NormalizeAdminSetting(
|
||||
"theme.catalog_url",
|
||||
|
||||
@@ -157,9 +157,10 @@ func (c MatcherConfig) TVSeriesRootQueueEnabled() bool {
|
||||
|
||||
// PlaybackConfig holds transcoding and playback settings.
|
||||
type PlaybackConfig struct {
|
||||
FFmpegPath string `yaml:"ffmpeg_path"`
|
||||
TranscodeDir string `yaml:"transcode_dir"`
|
||||
HWAccel string `yaml:"hw_accel"`
|
||||
FFmpegPath string `yaml:"ffmpeg_path"`
|
||||
TranscodeDir string `yaml:"transcode_dir"`
|
||||
SegmentRetentionSeconds int `yaml:"segment_retention_seconds"`
|
||||
HWAccel string `yaml:"hw_accel"`
|
||||
// HWDevice is the GPU render device for hardware transcodes. A single
|
||||
// path pins every GPU workload to that device; a comma-separated list
|
||||
// (e.g. "/dev/dri/renderD128,/dev/dri/renderD129") balances workloads
|
||||
@@ -479,6 +480,7 @@ func setDefaults() *configRaw {
|
||||
Playback: PlaybackConfig{
|
||||
FFmpegPath: "/usr/lib/jellyfin-ffmpeg/ffmpeg",
|
||||
TranscodeDir: DefaultTranscodeDir,
|
||||
SegmentRetentionSeconds: 600,
|
||||
HWAccel: "auto",
|
||||
ChapterThumbnailWorkers: 1,
|
||||
ChapterThumbnailExecution: "local",
|
||||
|
||||
@@ -317,6 +317,14 @@ func LoadFromDB(m map[string]string) (*Config, error) {
|
||||
// Playback
|
||||
cfg.Playback.FFmpegPath = stringOr(m, "playback.ffmpeg_path", "/usr/lib/jellyfin-ffmpeg/ffmpeg")
|
||||
cfg.Playback.TranscodeDir = stringOr(m, "playback.transcode_dir", DefaultTranscodeDir)
|
||||
segmentRetentionSeconds, err := intOr(m, playbackSegmentRetentionSettingKey, 600)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if segmentRetentionSeconds != 0 && (segmentRetentionSeconds < 120 || segmentRetentionSeconds > 86400) {
|
||||
return nil, fmt.Errorf("%s must be 0 or between 120 and 86400", playbackSegmentRetentionSettingKey)
|
||||
}
|
||||
cfg.Playback.SegmentRetentionSeconds = segmentRetentionSeconds
|
||||
cfg.Playback.HWAccel = stringOr(m, "playback.hw_accel", "auto")
|
||||
cfg.Playback.HWDevice = stringOr(m, "playback.hw_device", "")
|
||||
chapterThumbnailWorkers, err := intOr(m, "playback.chapter_thumbnail_workers", 1)
|
||||
|
||||
@@ -35,6 +35,17 @@ func TestLoadFromDBMetadataPresignExpiryRejectsInvalidDuration(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFromDBRejectsInvalidSegmentRetention(t *testing.T) {
|
||||
for _, value := range []string{"-1", "119", "86401"} {
|
||||
t.Run(value, func(t *testing.T) {
|
||||
_, err := LoadFromDB(map[string]string{playbackSegmentRetentionSettingKey: value})
|
||||
if err == nil || !strings.Contains(err.Error(), playbackSegmentRetentionSettingKey) {
|
||||
t.Fatalf("LoadFromDB() error = %v, want retention bounds error", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFromDBJellyfinWebEnabledDefaultsToTrue(t *testing.T) {
|
||||
cfg, err := LoadFromDB(map[string]string{})
|
||||
if err != nil {
|
||||
@@ -186,6 +197,16 @@ func TestYAMLToSettingsMapJellyfinCompatEnabledDefaultsToLegacyListener(t *testi
|
||||
}
|
||||
}
|
||||
|
||||
func TestYAMLToSettingsMapPreservesExplicitlyDisabledSegmentRetention(t *testing.T) {
|
||||
m := yamlSettingsMapFromString(t, `
|
||||
playback:
|
||||
segment_retention_seconds: 0
|
||||
`)
|
||||
if got := m[playbackSegmentRetentionSettingKey]; got != "0" {
|
||||
t.Fatalf("segment retention = %q, want explicit disable", got)
|
||||
}
|
||||
}
|
||||
|
||||
func yamlSettingsMapFromString(t *testing.T, body string) map[string]string {
|
||||
t.Helper()
|
||||
path := t.TempDir() + "/silo.yaml"
|
||||
|
||||
@@ -209,6 +209,7 @@ func YAMLToSettingsMap(path string) (map[string]string, error) {
|
||||
// Playback
|
||||
setIfNonEmpty(m, "playback.ffmpeg_path", raw.Playback.FFmpegPath)
|
||||
setIfNonEmpty(m, "playback.transcode_dir", raw.Playback.TranscodeDir)
|
||||
m[playbackSegmentRetentionSettingKey] = strconv.Itoa(raw.Playback.SegmentRetentionSeconds)
|
||||
setIfNonEmpty(m, "playback.hw_accel", raw.Playback.HWAccel)
|
||||
if raw.Playback.ChapterThumbnailWorkers != 0 {
|
||||
m["playback.chapter_thumbnail_workers"] = strconv.Itoa(raw.Playback.ChapterThumbnailWorkers)
|
||||
|
||||
@@ -184,6 +184,7 @@ type PlaybackHandler struct {
|
||||
FFmpegPath string
|
||||
HWAccel string
|
||||
TranscodeDir string
|
||||
SegmentRetentionSeconds int
|
||||
// tm is the shared transcode-session lifecycle (live map, reconstruct) — the
|
||||
// same type the native handler uses, so jellycompat gets the reconstruct cap
|
||||
// and node-affinity rule for free. The reconstruction recipe is carried in the
|
||||
@@ -268,36 +269,40 @@ func NewPlaybackHandler(
|
||||
transcodeDir := filepath.Join(os.TempDir(), "silo-transcode")
|
||||
ffmpegPath := ""
|
||||
hwAccel := ""
|
||||
segmentRetentionSeconds := 600
|
||||
if cfg != nil {
|
||||
if cfg.Playback.TranscodeDir != "" {
|
||||
transcodeDir = cfg.Playback.TranscodeDir
|
||||
}
|
||||
ffmpegPath = cfg.Playback.FFmpegPath
|
||||
hwAccel = cfg.Playback.HWAccel
|
||||
segmentRetentionSeconds = cfg.Playback.SegmentRetentionSeconds
|
||||
}
|
||||
|
||||
h := &PlaybackHandler{
|
||||
cfg: cfg,
|
||||
content: content,
|
||||
codec: codec,
|
||||
deviceProfiles: deviceProfiles,
|
||||
playbackStore: playbackStore,
|
||||
sessionMgr: sessionMgr,
|
||||
fileResolver: fileResolver,
|
||||
storeProvider: storeProvider,
|
||||
FFmpegPath: ffmpegPath,
|
||||
HWAccel: hwAccel,
|
||||
TranscodeDir: transcodeDir,
|
||||
tm: playback.NewTranscodeManager(),
|
||||
cfg: cfg,
|
||||
content: content,
|
||||
codec: codec,
|
||||
deviceProfiles: deviceProfiles,
|
||||
playbackStore: playbackStore,
|
||||
sessionMgr: sessionMgr,
|
||||
fileResolver: fileResolver,
|
||||
storeProvider: storeProvider,
|
||||
FFmpegPath: ffmpegPath,
|
||||
HWAccel: hwAccel,
|
||||
TranscodeDir: transcodeDir,
|
||||
SegmentRetentionSeconds: segmentRetentionSeconds,
|
||||
tm: playback.NewTranscodeManager(),
|
||||
}
|
||||
// Wire the shared transcode manager with closures so it reads the handler's
|
||||
// (late-set) JWTSecret lazily, matching the native handler.
|
||||
h.tm.JWTSecretFn = func() string { return h.JWTSecret }
|
||||
h.tm.Config = func() playback.TranscodeRuntimeConfig {
|
||||
return playback.TranscodeRuntimeConfig{
|
||||
TranscodeDir: h.TranscodeDir,
|
||||
FFmpegPath: h.FFmpegPath,
|
||||
HWAccel: h.HWAccel,
|
||||
TranscodeDir: h.TranscodeDir,
|
||||
FFmpegPath: h.FFmpegPath,
|
||||
HWAccel: h.HWAccel,
|
||||
SegmentRetentionSeconds: h.SegmentRetentionSeconds,
|
||||
}
|
||||
}
|
||||
if reg, ok := sessionMgr.(interface {
|
||||
|
||||
@@ -21,6 +21,7 @@ import (
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
|
||||
"github.com/Silo-Server/silo-server/internal/httpstream"
|
||||
"github.com/Silo-Server/silo-server/internal/models"
|
||||
"github.com/Silo-Server/silo-server/internal/nodepool"
|
||||
"github.com/Silo-Server/silo-server/internal/playback"
|
||||
@@ -415,8 +416,9 @@ func (h *PlaybackHandler) HandleHLSSegment(w http.ResponseWriter, r *http.Reques
|
||||
}
|
||||
}
|
||||
|
||||
segmentFile := name + "." + ext
|
||||
segmentPath, err := transcodeSession.GetSegment(segmentFile)
|
||||
segmentName := name + "." + ext
|
||||
downloadGeneration := transcodeSession.SegmentGeneration()
|
||||
segment, segmentInfo, err := transcodeSession.OpenSegment(segmentName)
|
||||
if err != nil && errors.Is(err, playback.ErrSegmentNotFound) {
|
||||
segNum, parseErr := playback.ParseSegmentNumber(name)
|
||||
if parseErr == nil {
|
||||
@@ -427,7 +429,7 @@ func (h *PlaybackHandler) HandleHLSSegment(w http.ResponseWriter, r *http.Reques
|
||||
lastProducedAgeMS = now.Sub(decision.Progress.LastProducedAt).Milliseconds()
|
||||
}
|
||||
slog.InfoContext(r.Context(), "transcode segment missing", "component", "jellycompat",
|
||||
"segment", segmentFile,
|
||||
"segment", segmentName,
|
||||
"requested_segment", segNum,
|
||||
"produced_head", decision.Progress.ProducedHead,
|
||||
"last_requested_segment", decision.Progress.LastRequestedSegment,
|
||||
@@ -441,7 +443,7 @@ func (h *PlaybackHandler) HandleHLSSegment(w http.ResponseWriter, r *http.Reques
|
||||
)
|
||||
if decision.Wait {
|
||||
slog.InfoContext(r.Context(), "transcode segment wait", "component", "jellycompat",
|
||||
"segment", segmentFile,
|
||||
"segment", segmentName,
|
||||
"requested_segment", segNum,
|
||||
"produced_head", decision.Progress.ProducedHead,
|
||||
"last_requested_segment", decision.Progress.LastRequestedSegment,
|
||||
@@ -453,10 +455,10 @@ func (h *PlaybackHandler) HandleHLSSegment(w http.ResponseWriter, r *http.Reques
|
||||
"session", playSession.UpstreamSessionID,
|
||||
"playback_session_id", playSession.UpstreamSessionID,
|
||||
)
|
||||
segmentPath, err = transcodeSession.WaitForSegment(segmentFile, decision.WaitTimeout)
|
||||
segment, segmentInfo, err = transcodeSession.WaitForOpenSegment(segmentName, decision.WaitTimeout)
|
||||
if err != nil && errors.Is(err, playback.ErrSegmentNotFound) {
|
||||
slog.InfoContext(r.Context(), "transcode segment wait timeout", "component", "jellycompat",
|
||||
"segment", segmentFile,
|
||||
"segment", segmentName,
|
||||
"requested_segment", segNum,
|
||||
"produced_head", decision.Progress.ProducedHead,
|
||||
"last_requested_segment", decision.Progress.LastRequestedSegment,
|
||||
@@ -476,7 +478,7 @@ func (h *PlaybackHandler) HandleHLSSegment(w http.ResponseWriter, r *http.Reques
|
||||
if seekErr != nil && !errors.Is(seekErr, playback.ErrManifestNotReady) {
|
||||
slog.ErrorContext(r.Context(), "resolve transcode seek target", "component", "jellycompat",
|
||||
"error", seekErr,
|
||||
"segment", segmentFile,
|
||||
"segment", segmentName,
|
||||
"play_session", playSessionID,
|
||||
"session", playSession.UpstreamSessionID,
|
||||
"playback_session_id", playSession.UpstreamSessionID,
|
||||
@@ -495,7 +497,7 @@ func (h *PlaybackHandler) HandleHLSSegment(w http.ResponseWriter, r *http.Reques
|
||||
|
||||
if ok {
|
||||
slog.InfoContext(r.Context(), "transcode seek restart", "component", "jellycompat",
|
||||
"segment", segmentFile,
|
||||
"segment", segmentName,
|
||||
"requested_segment", segNum,
|
||||
"produced_head", decision.Progress.ProducedHead,
|
||||
"last_requested_segment", decision.Progress.LastRequestedSegment,
|
||||
@@ -516,14 +518,14 @@ func (h *PlaybackHandler) HandleHLSSegment(w http.ResponseWriter, r *http.Reques
|
||||
seekSeconds,
|
||||
segNum,
|
||||
); restartErr == nil {
|
||||
segmentPath, err = transcodeSession.WaitForSegment(segmentFile, 30*time.Second)
|
||||
segment, segmentInfo, err = transcodeSession.WaitForOpenSegment(segmentName, 30*time.Second)
|
||||
}
|
||||
}
|
||||
}
|
||||
} else if transcodeSession.IsRunning() {
|
||||
// Non-numbered segment (e.g. init.mp4 for fMP4 HLS).
|
||||
// Wait briefly — the init segment is written almost immediately.
|
||||
segmentPath, err = transcodeSession.WaitForSegment(segmentFile, 10*time.Second)
|
||||
segment, segmentInfo, err = transcodeSession.WaitForOpenSegment(segmentName, 10*time.Second)
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
@@ -532,11 +534,17 @@ func (h *PlaybackHandler) HandleHLSSegment(w http.ResponseWriter, r *http.Reques
|
||||
return
|
||||
}
|
||||
|
||||
if segNum, parseErr := playback.ParseSegmentNumber(name); parseErr == nil {
|
||||
transcodeSession.ReportSegmentDownloaded(segNum)
|
||||
defer func() { _ = segment.Close() }()
|
||||
sw := httpstream.NewRollingDeadlineWriter(w)
|
||||
http.ServeContent(sw, r, segmentInfo.Name(), segmentInfo.ModTime(), segment)
|
||||
if r.Method == http.MethodGet &&
|
||||
sw.StatusCode() == http.StatusOK &&
|
||||
sw.BytesWritten() == segmentInfo.Size() &&
|
||||
sw.Outcome(r.Context()) == httpstream.OutcomeCompleted {
|
||||
if segNum, parseErr := playback.ParseSegmentNumber(name); parseErr == nil {
|
||||
transcodeSession.ReportSegmentDownloadedForGeneration(segNum, downloadGeneration)
|
||||
}
|
||||
}
|
||||
|
||||
http.ServeFile(w, r, segmentPath)
|
||||
}
|
||||
|
||||
// hlsSegmentErrorResponse maps a segment-retrieval error to a Jellyfin-faithful
|
||||
@@ -1655,6 +1663,7 @@ func (h *PlaybackHandler) ensureTranscodeSession(ctx context.Context, playSessio
|
||||
TotalDuration: float64(source.Version.Duration),
|
||||
FastStart: true,
|
||||
}
|
||||
opts.SegmentRetentionSeconds = h.SegmentRetentionSeconds
|
||||
if source.TranscodeAudio {
|
||||
opts.TargetCodecVideo = "copy"
|
||||
}
|
||||
|
||||
@@ -0,0 +1,119 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"time"
|
||||
)
|
||||
|
||||
const segmentPruneHysteresis = 5
|
||||
|
||||
// scheduleSegmentPruneLocked starts one asynchronous prune pass after the
|
||||
// client has advanced far enough to make useful work. The caller must hold
|
||||
// s.mu. Segment retention is measured in media time so custom HLS segment
|
||||
// durations retain the same back-seek window.
|
||||
func (s *TranscodeSession) scheduleSegmentPruneLocked() {
|
||||
retentionSeconds := s.opts.SegmentRetentionSeconds
|
||||
if retentionSeconds <= 0 || s.segmentPruneRunning || s.restarting {
|
||||
return
|
||||
}
|
||||
segmentDuration := s.opts.SegmentDuration
|
||||
if segmentDuration <= 0 {
|
||||
segmentDuration = defaultSegmentDuration
|
||||
}
|
||||
retainedSegments := (retentionSeconds + segmentDuration - 1) / segmentDuration
|
||||
floor := s.lastRequestedSegment - retainedSegments
|
||||
if floor <= s.opts.StartSegmentNumber || floor-s.lastPruneFloor < segmentPruneHysteresis {
|
||||
return
|
||||
}
|
||||
|
||||
s.segmentPruneRunning = true
|
||||
generation := s.segmentGeneration
|
||||
go s.pruneDownloadedSegments(generation, floor)
|
||||
}
|
||||
|
||||
// pruneDownloadedSegments removes completed media files strictly behind floor.
|
||||
// FFmpeg's manifest and init files are never candidates, and the current
|
||||
// process's startup window remains present so real-manifest reloads cannot get
|
||||
// stuck behind startupFilesReady after cleanup begins.
|
||||
func (s *TranscodeSession) pruneDownloadedSegments(generation uint64, floor int) {
|
||||
started := time.Now()
|
||||
entries, err := os.ReadDir(s.outputDir)
|
||||
if err != nil {
|
||||
s.finishSegmentPrune(generation, floor)
|
||||
if !errors.Is(err, os.ErrNotExist) {
|
||||
slog.Warn("read transcode segments for pruning", "component", "playback", "error", err, "session", s.opts.SessionID, "playback_session_id", s.opts.SessionID)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
opts := s.opts
|
||||
s.mu.Unlock()
|
||||
segmentDuration := opts.SegmentDuration
|
||||
if segmentDuration <= 0 {
|
||||
segmentDuration = defaultSegmentDuration
|
||||
}
|
||||
freshAfter := time.Now().Add(-(2*time.Duration(segmentDuration)*time.Second + 30*time.Second))
|
||||
startupEnd := opts.StartSegmentNumber + startupSegmentRequirement(opts)
|
||||
|
||||
removed := 0
|
||||
var freedBytes int64
|
||||
for _, entry := range entries {
|
||||
segment, parseErr := ParseSegmentNumber(entry.Name())
|
||||
if parseErr != nil || segment >= floor || (segment >= opts.StartSegmentNumber && segment < startupEnd) {
|
||||
continue
|
||||
}
|
||||
info, infoErr := entry.Info()
|
||||
if infoErr != nil || info.ModTime().After(freshAfter) {
|
||||
continue
|
||||
}
|
||||
|
||||
// Serialize the unlink with restart's generation change. Once restart
|
||||
// advances the generation, this pass cannot delete newly generated files.
|
||||
s.mu.Lock()
|
||||
if generation != s.segmentGeneration || s.restarting {
|
||||
s.mu.Unlock()
|
||||
s.finishSegmentPrune(generation, floor)
|
||||
return
|
||||
}
|
||||
removeErr := os.Remove(filepath.Join(s.outputDir, entry.Name()))
|
||||
s.mu.Unlock()
|
||||
if removeErr != nil {
|
||||
if !errors.Is(removeErr, os.ErrNotExist) {
|
||||
slog.Warn("remove downloaded transcode segment", "component", "playback", "error", removeErr, "segment", segment, "session", opts.SessionID, "playback_session_id", opts.SessionID)
|
||||
}
|
||||
continue
|
||||
}
|
||||
removed++
|
||||
freedBytes += info.Size()
|
||||
}
|
||||
|
||||
s.finishSegmentPrune(generation, floor)
|
||||
if removed > 0 {
|
||||
slog.Info("pruned downloaded transcode segments",
|
||||
"component", "playback",
|
||||
"count", removed,
|
||||
"freed_bytes", freedBytes,
|
||||
"floor_segment", floor,
|
||||
"duration_ms", time.Since(started).Milliseconds(),
|
||||
"session", opts.SessionID,
|
||||
"playback_session_id", opts.SessionID,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *TranscodeSession) finishSegmentPrune(generation uint64, floor int) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if generation != s.segmentGeneration {
|
||||
return
|
||||
}
|
||||
if floor > s.lastPruneFloor {
|
||||
s.lastPruneFloor = floor
|
||||
}
|
||||
s.segmentPruneRunning = false
|
||||
s.scheduleSegmentPruneLocked()
|
||||
}
|
||||
@@ -0,0 +1,145 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestReportSegmentDownloadedPrunesOnlyExpiredBackBuffer(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
old := time.Now().Add(-time.Minute)
|
||||
for segment := 0; segment <= 15; segment++ {
|
||||
writePrunerTestFile(t, filepath.Join(dir, segmentFilename(segment, TranscodeOpts{})), []byte("segment"), old)
|
||||
}
|
||||
writePrunerTestFile(t, filepath.Join(dir, "init.mp4"), []byte("init"), old)
|
||||
writePrunerTestFile(t, filepath.Join(dir, "stream.m3u8"), []byte("manifest"), old)
|
||||
writePrunerTestFile(t, filepath.Join(dir, "seg_00004.ts.tmp"), []byte("temporary"), old)
|
||||
writePrunerTestFile(t, filepath.Join(dir, "notes.txt"), []byte("other"), old)
|
||||
// A newly-created file below the floor may belong to a replacement ffmpeg
|
||||
// process and must survive this pass.
|
||||
writePrunerTestFile(t, filepath.Join(dir, segmentFilename(4, TranscodeOpts{})), []byte("fresh"), time.Now())
|
||||
|
||||
session := &TranscodeSession{
|
||||
opts: TranscodeOpts{
|
||||
SessionID: "prune-test",
|
||||
SegmentDuration: 2,
|
||||
SegmentRetentionSeconds: 10,
|
||||
},
|
||||
outputDir: dir,
|
||||
lastPruneFloor: -1,
|
||||
lastRequestedSegment: 0,
|
||||
}
|
||||
session.ReportSegmentDownloaded(15)
|
||||
waitForPrunePass(t, session)
|
||||
|
||||
for _, segment := range []int{0, 1, 2, 4, 10, 11, 12, 13, 14, 15} {
|
||||
assertPrunerFileExists(t, filepath.Join(dir, segmentFilename(segment, TranscodeOpts{})))
|
||||
}
|
||||
for _, segment := range []int{3, 5, 6, 7, 8, 9} {
|
||||
assertPrunerFileMissing(t, filepath.Join(dir, segmentFilename(segment, TranscodeOpts{})))
|
||||
}
|
||||
for _, name := range []string{"init.mp4", "stream.m3u8", "seg_00004.ts.tmp", "notes.txt"} {
|
||||
assertPrunerFileExists(t, filepath.Join(dir, name))
|
||||
}
|
||||
}
|
||||
|
||||
func TestReportSegmentDownloadedDisabledKeepsSegments(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
path := filepath.Join(dir, segmentFilename(3, TranscodeOpts{}))
|
||||
writePrunerTestFile(t, path, []byte("segment"), time.Now().Add(-time.Minute))
|
||||
session := &TranscodeSession{
|
||||
opts: TranscodeOpts{SegmentDuration: 2},
|
||||
outputDir: dir,
|
||||
lastPruneFloor: -1,
|
||||
lastRequestedSegment: 0,
|
||||
}
|
||||
|
||||
session.ReportSegmentDownloaded(500)
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
assertPrunerFileExists(t, path)
|
||||
}
|
||||
|
||||
func TestOldGenerationDownloadDoesNotAdvanceRestartedSession(t *testing.T) {
|
||||
session := &TranscodeSession{
|
||||
opts: TranscodeOpts{SegmentRetentionSeconds: 600},
|
||||
lastRequestedSegment: 100,
|
||||
}
|
||||
oldGeneration := session.SegmentGeneration()
|
||||
|
||||
session.mu.Lock()
|
||||
session.segmentGeneration++
|
||||
session.lastRequestedSegment = 20
|
||||
session.mu.Unlock()
|
||||
|
||||
session.ReportSegmentDownloadedForGeneration(101, oldGeneration)
|
||||
if got := session.LastRequestedSegment(); got != 20 {
|
||||
t.Fatalf("last requested segment = %d, want restarted position 20", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenSegmentDescriptorSurvivesUnlink(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
path := filepath.Join(dir, "seg_00001.ts")
|
||||
want := []byte("complete segment")
|
||||
writePrunerTestFile(t, path, want, time.Now())
|
||||
session := &TranscodeSession{outputDir: dir}
|
||||
|
||||
segment, _, err := session.OpenSegment("seg_00001.ts")
|
||||
if err != nil {
|
||||
t.Fatalf("OpenSegment: %v", err)
|
||||
}
|
||||
defer func() { _ = segment.Close() }()
|
||||
if err := os.Remove(path); err != nil {
|
||||
t.Fatalf("remove opened segment: %v", err)
|
||||
}
|
||||
got, err := io.ReadAll(segment)
|
||||
if err != nil {
|
||||
t.Fatalf("read opened segment: %v", err)
|
||||
}
|
||||
if string(got) != string(want) {
|
||||
t.Fatalf("opened segment = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func writePrunerTestFile(t *testing.T, path string, contents []byte, modTime time.Time) {
|
||||
t.Helper()
|
||||
if err := os.WriteFile(path, contents, 0o600); err != nil {
|
||||
t.Fatalf("write %s: %v", filepath.Base(path), err)
|
||||
}
|
||||
if err := os.Chtimes(path, modTime, modTime); err != nil {
|
||||
t.Fatalf("set times for %s: %v", filepath.Base(path), err)
|
||||
}
|
||||
}
|
||||
|
||||
func waitForPrunePass(t *testing.T, session *TranscodeSession) {
|
||||
t.Helper()
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
session.mu.Lock()
|
||||
running := session.segmentPruneRunning
|
||||
floor := session.lastPruneFloor
|
||||
session.mu.Unlock()
|
||||
if !running && floor >= 10 {
|
||||
return
|
||||
}
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
}
|
||||
t.Fatal("segment prune pass did not finish")
|
||||
}
|
||||
|
||||
func assertPrunerFileExists(t *testing.T, path string) {
|
||||
t.Helper()
|
||||
if _, err := os.Stat(path); err != nil {
|
||||
t.Fatalf("expected %s to exist: %v", filepath.Base(path), err)
|
||||
}
|
||||
}
|
||||
|
||||
func assertPrunerFileMissing(t *testing.T, path string) {
|
||||
t.Helper()
|
||||
if _, err := os.Stat(path); !os.IsNotExist(err) {
|
||||
t.Fatalf("expected %s to be removed, stat error = %v", filepath.Base(path), err)
|
||||
}
|
||||
}
|
||||
+113
-11
@@ -45,17 +45,18 @@ type TranscodeOpts struct {
|
||||
StreamOriginSeconds float64
|
||||
// CopySeekAnchorResolved distinguishes a valid zero-second origin from
|
||||
// older/shared recipes that never resolved a copy seek anchor.
|
||||
CopySeekAnchorResolved bool
|
||||
TargetResolution string // e.g., 1080p, 720p
|
||||
TargetCodecVideo string // e.g., h264 (or hevc if allowed)
|
||||
TargetCodecAudio string // e.g., aac
|
||||
SegmentDuration int // seconds, default 6
|
||||
StartSegmentNumber int // -hls_segment_start_number, default 0
|
||||
FFmpegPath string // optional explicit ffmpeg binary path
|
||||
HWAccel string // auto, qsv, vaapi, nvenc, none
|
||||
HWDevice string // e.g., /dev/dri/renderD128 (default if empty)
|
||||
SubtitleTrackIndex int // -1 = no subtitles
|
||||
SubtitleBurnIn bool
|
||||
CopySeekAnchorResolved bool
|
||||
TargetResolution string // e.g., 1080p, 720p
|
||||
TargetCodecVideo string // e.g., h264 (or hevc if allowed)
|
||||
TargetCodecAudio string // e.g., aac
|
||||
SegmentDuration int // seconds, default 6
|
||||
SegmentRetentionSeconds int // downloaded media retained behind the client; 0 disables pruning
|
||||
StartSegmentNumber int // -hls_segment_start_number, default 0
|
||||
FFmpegPath string // optional explicit ffmpeg binary path
|
||||
HWAccel string // auto, qsv, vaapi, nvenc, none
|
||||
HWDevice string // e.g., /dev/dri/renderD128 (default if empty)
|
||||
SubtitleTrackIndex int // -1 = no subtitles
|
||||
SubtitleBurnIn bool
|
||||
// SubtitleCodec is the probed codec of the burn-in track (e.g. "subrip",
|
||||
// "hdmv_pgs_subtitle"). Bitmap codecs (PGS/DVD/DVB) select the overlay
|
||||
// filter_complex pipeline; text codecs use the libass subtitles filter.
|
||||
@@ -92,6 +93,9 @@ type TranscodeSession struct {
|
||||
done chan struct{} // closed when the monitor goroutine finishes
|
||||
stdinPipe io.WriteCloser
|
||||
lastRequestedSegment int
|
||||
lastPruneFloor int
|
||||
segmentPruneRunning bool
|
||||
segmentGeneration uint64
|
||||
throttler *TranscodeThrottler
|
||||
stderrLinesLogged int
|
||||
stderrBytesLogged int
|
||||
@@ -202,6 +206,7 @@ func StartTranscode(ctx context.Context, opts TranscodeOpts) (*TranscodeSession,
|
||||
done: make(chan struct{}),
|
||||
stderr: newBoundedTailBuffer(stderrTailMaxBytes),
|
||||
lastRequestedSegment: opts.StartSegmentNumber,
|
||||
lastPruneFloor: opts.StartSegmentNumber - 1,
|
||||
reserveHWDeviceOnRestart: reserveHWDeviceOnRestart,
|
||||
}
|
||||
|
||||
@@ -1620,6 +1625,31 @@ func (s *TranscodeSession) GetSegment(name string) (string, error) {
|
||||
return segPath, nil
|
||||
}
|
||||
|
||||
// OpenSegment opens a completed segment for serving. Keeping the descriptor
|
||||
// open while the response is written makes it a filesystem-backed lease: a
|
||||
// concurrent prune may unlink the directory entry, but the response can still
|
||||
// read the complete file on POSIX filesystems.
|
||||
func (s *TranscodeSession) OpenSegment(name string) (*os.File, os.FileInfo, error) {
|
||||
clean := filepath.Base(name)
|
||||
segment, err := os.Open(filepath.Join(s.outputDir, clean))
|
||||
if err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
return nil, nil, ErrSegmentNotFound
|
||||
}
|
||||
return nil, nil, fmt.Errorf("open segment: %w", err)
|
||||
}
|
||||
info, err := segment.Stat()
|
||||
if err != nil {
|
||||
_ = segment.Close()
|
||||
return nil, nil, fmt.Errorf("stat segment: %w", err)
|
||||
}
|
||||
if info.Size() <= 0 {
|
||||
_ = segment.Close()
|
||||
return nil, nil, ErrSegmentNotFound
|
||||
}
|
||||
return segment, info, nil
|
||||
}
|
||||
|
||||
// Close terminates the ffmpeg process and removes the temporary output directory.
|
||||
func (s *TranscodeSession) Close() error {
|
||||
return s.shutdown(true)
|
||||
@@ -1769,6 +1799,10 @@ func (s *TranscodeSession) restart(
|
||||
return nil
|
||||
}
|
||||
s.restarting = true
|
||||
s.segmentGeneration++
|
||||
s.segmentPruneRunning = false
|
||||
s.lastRequestedSegment = startSegment
|
||||
s.lastPruneFloor = startSegment - 1
|
||||
cancelCurrent := s.cancel
|
||||
done := s.done
|
||||
s.mu.Unlock()
|
||||
@@ -1928,6 +1962,49 @@ func (s *TranscodeSession) WaitForSegment(name string, timeout time.Duration) (s
|
||||
}
|
||||
}
|
||||
|
||||
// WaitForOpenSegment is the opened-file counterpart to WaitForSegment. It
|
||||
// closes the stat-to-open race for HTTP serving while preserving the same
|
||||
// restart and ffmpeg-exit error contract.
|
||||
func (s *TranscodeSession) WaitForOpenSegment(name string, timeout time.Duration) (*os.File, os.FileInfo, error) {
|
||||
deadline := time.After(timeout)
|
||||
for {
|
||||
segment, info, err := s.OpenSegment(name)
|
||||
if err == nil {
|
||||
return segment, info, nil
|
||||
}
|
||||
if !errors.Is(err, ErrSegmentNotFound) {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
running := s.running
|
||||
restarting := s.restarting
|
||||
waitErr := s.waitErr
|
||||
s.mu.Unlock()
|
||||
|
||||
if restarting {
|
||||
select {
|
||||
case <-deadline:
|
||||
return nil, nil, ErrSegmentNotFound
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
continue
|
||||
}
|
||||
}
|
||||
if !running && waitErr != nil {
|
||||
return nil, nil, fmt.Errorf("%w: %w", ErrTranscodeFailed, waitErr)
|
||||
}
|
||||
if !running {
|
||||
return nil, nil, ErrSegmentNotFound
|
||||
}
|
||||
|
||||
select {
|
||||
case <-deadline:
|
||||
return nil, nil, ErrSegmentNotFound
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// RewriteManifestPaths prefixes relative segment references in an HLS manifest
|
||||
// with segPrefix (e.g. "segment/") and optionally appends rawQuery as a query
|
||||
// string. This ensures the HLS player's segment requests match server routes
|
||||
@@ -2175,9 +2252,34 @@ func (s *TranscodeSession) RestartSeekTarget(segNum int) (float64, bool, error)
|
||||
func (s *TranscodeSession) ReportSegmentDownloaded(segNum int) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.reportSegmentDownloadedLocked(segNum)
|
||||
}
|
||||
|
||||
// SegmentGeneration identifies the current ffmpeg timeline. Callers that hold
|
||||
// an open segment across a restart use it to prevent the old response from
|
||||
// advancing the replacement process's download position.
|
||||
func (s *TranscodeSession) SegmentGeneration() uint64 {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.segmentGeneration
|
||||
}
|
||||
|
||||
// ReportSegmentDownloadedForGeneration records completion only if the served
|
||||
// file belongs to the current ffmpeg timeline.
|
||||
func (s *TranscodeSession) ReportSegmentDownloadedForGeneration(segNum int, generation uint64) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if generation != s.segmentGeneration {
|
||||
return
|
||||
}
|
||||
s.reportSegmentDownloadedLocked(segNum)
|
||||
}
|
||||
|
||||
func (s *TranscodeSession) reportSegmentDownloadedLocked(segNum int) {
|
||||
if segNum > s.lastRequestedSegment {
|
||||
s.lastRequestedSegment = segNum
|
||||
}
|
||||
s.scheduleSegmentPruneLocked()
|
||||
}
|
||||
|
||||
// LastRequestedSegment returns the highest segment number downloaded by the client.
|
||||
|
||||
@@ -20,10 +20,11 @@ import (
|
||||
// internal/playback does not import internal/config (avoiding an import cycle);
|
||||
// each embedding handler adapts its own config snapshot into this shape.
|
||||
type TranscodeRuntimeConfig struct {
|
||||
TranscodeDir string
|
||||
FFmpegPath string
|
||||
HWAccel string
|
||||
HWDevice string
|
||||
TranscodeDir string
|
||||
FFmpegPath string
|
||||
HWAccel string
|
||||
HWDevice string
|
||||
SegmentRetentionSeconds int
|
||||
}
|
||||
|
||||
// sessionReconstructor is the SessionManager capability used to re-register a
|
||||
@@ -577,6 +578,7 @@ func (m *TranscodeManager) doReconstructTranscode(ctx context.Context, sessionID
|
||||
// operator config change applies to reconstructed sessions too.
|
||||
opts.HWAccel = cfg.HWAccel
|
||||
opts.HWDevice = cfg.HWDevice
|
||||
opts.SegmentRetentionSeconds = cfg.SegmentRetentionSeconds
|
||||
|
||||
// Resume near the segment the client is actually requesting. The card records
|
||||
// the original start; if the client has played past it, spawning ffmpeg at the
|
||||
|
||||
@@ -19,6 +19,7 @@ import (
|
||||
"golang.org/x/sync/singleflight"
|
||||
|
||||
"github.com/Silo-Server/silo-server/internal/chapterthumbs"
|
||||
"github.com/Silo-Server/silo-server/internal/httpstream"
|
||||
"github.com/Silo-Server/silo-server/internal/nodeconfig"
|
||||
"github.com/Silo-Server/silo-server/internal/nodesessions"
|
||||
"github.com/Silo-Server/silo-server/internal/playback"
|
||||
@@ -476,23 +477,24 @@ func (s *Server) handleStart(w http.ResponseWriter, r *http.Request) {
|
||||
outputDir := filepath.Join(cfg.Playback.TranscodeDir, req.SessionID)
|
||||
|
||||
opts := playback.TranscodeOpts{
|
||||
InputPath: req.InputPath,
|
||||
OutputDir: outputDir,
|
||||
SessionID: req.SessionID,
|
||||
SourceVideoCodec: req.SourceVideoCodec,
|
||||
VideoBitstreamFilter: req.VideoBitstreamFilter,
|
||||
SeekSeconds: req.SeekSeconds,
|
||||
StreamOriginSeconds: req.StreamOriginSeconds,
|
||||
CopySeekAnchorResolved: req.CopySeekAnchorResolved,
|
||||
StartSegmentNumber: req.StartSegmentNumber,
|
||||
TargetResolution: req.TargetResolution,
|
||||
TargetCodecVideo: req.TargetCodecVideo,
|
||||
TargetCodecAudio: req.TargetCodecAudio,
|
||||
TargetAudioChannels: req.TargetAudioChannels,
|
||||
TargetBitrateKbps: req.TargetBitrateKbps,
|
||||
SegmentDuration: req.SegmentDuration,
|
||||
FFmpegPath: cfg.Playback.FFmpegPath,
|
||||
HWAccel: req.HWAccel,
|
||||
InputPath: req.InputPath,
|
||||
OutputDir: outputDir,
|
||||
SessionID: req.SessionID,
|
||||
SourceVideoCodec: req.SourceVideoCodec,
|
||||
VideoBitstreamFilter: req.VideoBitstreamFilter,
|
||||
SeekSeconds: req.SeekSeconds,
|
||||
StreamOriginSeconds: req.StreamOriginSeconds,
|
||||
CopySeekAnchorResolved: req.CopySeekAnchorResolved,
|
||||
StartSegmentNumber: req.StartSegmentNumber,
|
||||
TargetResolution: req.TargetResolution,
|
||||
TargetCodecVideo: req.TargetCodecVideo,
|
||||
TargetCodecAudio: req.TargetCodecAudio,
|
||||
TargetAudioChannels: req.TargetAudioChannels,
|
||||
TargetBitrateKbps: req.TargetBitrateKbps,
|
||||
SegmentDuration: req.SegmentDuration,
|
||||
SegmentRetentionSeconds: cfg.Playback.SegmentRetentionSeconds,
|
||||
FFmpegPath: cfg.Playback.FFmpegPath,
|
||||
HWAccel: req.HWAccel,
|
||||
// This node's configured device (or device list — StartTranscode
|
||||
// resolves it to one GPU), matching what reconstruction uses so fresh
|
||||
// and reconstructed sessions balance identically.
|
||||
@@ -697,6 +699,7 @@ func (s *Server) spawnReconstruct(r *http.Request, sessionID string, requestedSe
|
||||
// rebuild. Run as a transcode node, not integrated (card.TranscodeOpts defaults).
|
||||
opts.HWAccel = cfg.Playback.HWAccel
|
||||
opts.HWDevice = cfg.Playback.HWDevice
|
||||
opts.SegmentRetentionSeconds = cfg.Playback.SegmentRetentionSeconds
|
||||
opts.NodeType = "transcode"
|
||||
opts.ExecutionMode = "transcode_node"
|
||||
|
||||
@@ -891,7 +894,8 @@ func (s *Server) handleSegment(w http.ResponseWriter, r *http.Request) {
|
||||
s.touchSession(sessionID)
|
||||
}
|
||||
|
||||
segPath, err := session.GetSegment(name)
|
||||
downloadGeneration := session.SegmentGeneration()
|
||||
segment, segmentInfo, err := session.OpenSegment(name)
|
||||
if err != nil && err == playback.ErrSegmentNotFound {
|
||||
segNum, parseErr := playback.ParseSegmentNumber(name)
|
||||
if parseErr == nil {
|
||||
@@ -926,7 +930,7 @@ func (s *Server) handleSegment(w http.ResponseWriter, r *http.Request) {
|
||||
"session", sessionID,
|
||||
"playback_session_id", sessionID,
|
||||
)
|
||||
segPath, err = session.WaitForSegment(name, decision.WaitTimeout)
|
||||
segment, segmentInfo, err = session.WaitForOpenSegment(name, decision.WaitTimeout)
|
||||
if err != nil && err == playback.ErrSegmentNotFound {
|
||||
slog.InfoContext(r.Context(), "transcode segment wait timeout", "component", "transcodenode",
|
||||
"segment", name,
|
||||
@@ -971,7 +975,7 @@ func (s *Server) handleSegment(w http.ResponseWriter, r *http.Request) {
|
||||
seekSeconds,
|
||||
segNum,
|
||||
); restartErr == nil {
|
||||
segPath, err = session.WaitForSegment(name, 30*time.Second)
|
||||
segment, segmentInfo, err = session.WaitForOpenSegment(name, 30*time.Second)
|
||||
}
|
||||
}
|
||||
if !ok && session.IsCopyVideo() {
|
||||
@@ -981,7 +985,7 @@ func (s *Server) handleSegment(w http.ResponseWriter, r *http.Request) {
|
||||
} else if session.IsRunning() {
|
||||
// Non-numbered segment (e.g., init.mp4 for fMP4 HLS).
|
||||
// Wait briefly — the init segment is written almost immediately.
|
||||
segPath, err = session.WaitForSegment(name, 10*time.Second)
|
||||
segment, segmentInfo, err = session.WaitForOpenSegment(name, 10*time.Second)
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
@@ -991,7 +995,17 @@ func (s *Server) handleSegment(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
w.Header().Set("Cache-Control", "no-store, max-age=0")
|
||||
w.Header().Set("Pragma", "no-cache")
|
||||
http.ServeFile(w, r, segPath)
|
||||
defer func() { _ = segment.Close() }()
|
||||
sw := httpstream.NewRollingDeadlineWriter(w)
|
||||
http.ServeContent(sw, r, segmentInfo.Name(), segmentInfo.ModTime(), segment)
|
||||
if r.Method == http.MethodGet &&
|
||||
sw.StatusCode() == http.StatusOK &&
|
||||
sw.BytesWritten() == segmentInfo.Size() &&
|
||||
sw.Outcome(r.Context()) == httpstream.OutcomeCompleted {
|
||||
if segmentNumber, parseErr := playback.ParseSegmentNumber(name); parseErr == nil {
|
||||
session.ReportSegmentDownloadedForGeneration(segmentNumber, downloadGeneration)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Server) handleForceReload(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
@@ -16,6 +16,7 @@ import {
|
||||
const KEYS = [
|
||||
"playback.ffmpeg_path",
|
||||
"playback.transcode_dir",
|
||||
"playback.segment_retention_seconds",
|
||||
"playback.hw_accel",
|
||||
"playback.hw_device",
|
||||
"playback.transcode_enabled",
|
||||
@@ -196,6 +197,13 @@ export default function PlaybackSettings() {
|
||||
</FieldGroup>
|
||||
|
||||
<FieldGroup label="Segments">
|
||||
<SettingField
|
||||
label="Transcode Back Buffer (seconds)"
|
||||
type="number"
|
||||
hint="Keeps this much already-downloaded media for instant backward seeking, then reclaims older transcode segments. Use 0 to disable cleanup; enabled values must be at least 120 seconds. Pair with transcode throttling to bound both behind- and ahead-of-client disk usage."
|
||||
value={form.getValue("playback.segment_retention_seconds")}
|
||||
onChange={(v) => form.setValue("playback.segment_retention_seconds", v)}
|
||||
/>
|
||||
<SettingField
|
||||
label="Chapter Thumbnail Workers"
|
||||
type="number"
|
||||
|
||||
Reference in New Issue
Block a user