From 6a3fcef165706f4bcc21140f40b82fbb65ba040c Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Sat, 8 Aug 2026 15:23:16 -0400 Subject: [PATCH] fix(playback): prune downloaded transcode segments --- internal/api/handlers/playback.go | 97 ++++++------ internal/api/handlers/playback_v3.go | 2 +- internal/config/admin_settings.go | 16 +- internal/config/admin_settings_test.go | 15 ++ internal/config/config.go | 8 +- internal/config/db_loader.go | 8 + internal/config/db_loader_test.go | 21 +++ internal/config/yaml_import.go | 1 + internal/jellycompat/handlers_playback.go | 35 +++-- internal/jellycompat/streams.go | 37 +++-- internal/playback/segment_pruner.go | 119 ++++++++++++++ internal/playback/segment_pruner_test.go | 145 ++++++++++++++++++ internal/playback/transcode.go | 124 +++++++++++++-- internal/playback/transcode_manager.go | 10 +- internal/transcodenode/server.go | 58 ++++--- .../pages/admin-settings/PlaybackSettings.tsx | 8 + 16 files changed, 590 insertions(+), 114 deletions(-) create mode 100644 internal/playback/segment_pruner.go create mode 100644 internal/playback/segment_pruner_test.go diff --git a/internal/api/handlers/playback.go b/internal/api/handlers/playback.go index c22a4103..78c99192 100644 --- a/internal/api/handlers/playback.go +++ b/internal/api/handlers/playback.go @@ -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 diff --git a/internal/api/handlers/playback_v3.go b/internal/api/handlers/playback_v3.go index 627a07b8..db03e17f 100644 --- a/internal/api/handlers/playback_v3.go +++ b/internal/api/handlers/playback_v3.go @@ -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} diff --git a/internal/config/admin_settings.go b/internal/config/admin_settings.go index a3f1aae1..236f20ca 100644 --- a/internal/config/admin_settings.go +++ b/internal/config/admin_settings.go @@ -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": diff --git a/internal/config/admin_settings_test.go b/internal/config/admin_settings_test.go index 6351534e..ad5918a4 100644 --- a/internal/config/admin_settings_test.go +++ b/internal/config/admin_settings_test.go @@ -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", diff --git a/internal/config/config.go b/internal/config/config.go index 82c30a32..2d362173 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -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", diff --git a/internal/config/db_loader.go b/internal/config/db_loader.go index fcef840e..4462a48a 100644 --- a/internal/config/db_loader.go +++ b/internal/config/db_loader.go @@ -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) diff --git a/internal/config/db_loader_test.go b/internal/config/db_loader_test.go index 70a414a7..f24bfead 100644 --- a/internal/config/db_loader_test.go +++ b/internal/config/db_loader_test.go @@ -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" diff --git a/internal/config/yaml_import.go b/internal/config/yaml_import.go index fbd903da..73d2249a 100644 --- a/internal/config/yaml_import.go +++ b/internal/config/yaml_import.go @@ -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) diff --git a/internal/jellycompat/handlers_playback.go b/internal/jellycompat/handlers_playback.go index dd7c2bde..e5d4b85e 100644 --- a/internal/jellycompat/handlers_playback.go +++ b/internal/jellycompat/handlers_playback.go @@ -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 { diff --git a/internal/jellycompat/streams.go b/internal/jellycompat/streams.go index cc240a6e..c9e2fd96 100644 --- a/internal/jellycompat/streams.go +++ b/internal/jellycompat/streams.go @@ -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" } diff --git a/internal/playback/segment_pruner.go b/internal/playback/segment_pruner.go new file mode 100644 index 00000000..c75a8133 --- /dev/null +++ b/internal/playback/segment_pruner.go @@ -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() +} diff --git a/internal/playback/segment_pruner_test.go b/internal/playback/segment_pruner_test.go new file mode 100644 index 00000000..6a27cabf --- /dev/null +++ b/internal/playback/segment_pruner_test.go @@ -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) + } +} diff --git a/internal/playback/transcode.go b/internal/playback/transcode.go index c2bee799..6c5abbef 100644 --- a/internal/playback/transcode.go +++ b/internal/playback/transcode.go @@ -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. diff --git a/internal/playback/transcode_manager.go b/internal/playback/transcode_manager.go index 308713a1..08747dec 100644 --- a/internal/playback/transcode_manager.go +++ b/internal/playback/transcode_manager.go @@ -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 diff --git a/internal/transcodenode/server.go b/internal/transcodenode/server.go index 57c397ef..f3d04b7c 100644 --- a/internal/transcodenode/server.go +++ b/internal/transcodenode/server.go @@ -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) { diff --git a/web/src/pages/admin-settings/PlaybackSettings.tsx b/web/src/pages/admin-settings/PlaybackSettings.tsx index 0df24b52..a1f9e027 100644 --- a/web/src/pages/admin-settings/PlaybackSettings.tsx +++ b/web/src/pages/admin-settings/PlaybackSettings.tsx @@ -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() { + form.setValue("playback.segment_retention_seconds", v)} + />