fix(playback): open the listener before sweeping stale transcode dirs (#413)
* fix(playback): run orphaned-transcode cleanup in the background at startup The native and Jellyfin-compat routers swept stale per-session transcode dirs synchronously during NewRouter, before the listener bound. On a slow network filesystem this blocked startup for 80+s (64 leftover dirs on the last deploy), so restart-reconnect clients were turned away and the health check reported the server unhealthy the whole time. Move both sweeps into a background goroutine (StartBackgroundOrphanCleanup) so the listener comes up immediately and the cleanup runs concurrently. The delete logic is unchanged: same active-session snapshot and MaxTokenTTL age-sparing, only later. A package-level mutex serializes concurrent sweeps of the shared transcode root so the two background sweeps can't race on os.RemoveAll. Part of #412 * fix(transcode): background the node boot-time transcode-dir sweep A dedicated transcode node swept leftover transcode dirs synchronously in NewServer, before startStandaloneServer bound its listener. On a slow network filesystem that delete blocked the node from coming online at boot, the same startup-stall class as the main server. Move the sweep into the shared StartBackgroundOrphanCleanup goroutine so the node's listener binds immediately. Backgrounding required an age guard: the sweep previously ran as a full wipe (minAge=0) with an empty active-set, which was only safe because it completed before any request could arrive. Run concurrently that would race a token-carried reconstruct writing into TranscodeDir/<sessionID>, deleting segments a fresh ffmpeg is producing. Passing MaxTokenTTL spares any dir younger than the max token lifetime — exactly the ones a still-valid reconnect could reconstruct — while dirs older than any surviving token (never reconstructable) are still reclaimed. Part of #412 * feat(playback): reclaim orphaned transcode dirs periodically, not just at boot The orphaned-transcode sweep only ran at startup on both the central server and transcode nodes, so it only ever reclaimed dirs left by an ungraceful prior shutdown. During a long uptime the in-memory session reapers delete the dirs of sessions they still track, but a dir whose owning session was dropped without its RemoveAll succeeding becomes an "untracked orphan" with no runtime GC — on a box that runs for weeks these accumulate until the next restart. Add StartPeriodicOrphanCleanup: an immediate background sweep followed by an hourly re-run bound to a lifecycle context. Wire it on all three surfaces — native API and Jellyfin-compat (via deps.AppContext) and the transcode node (via a new Server.StartOrphanSweeper(appCtx), replacing its boot-only sweep). When no context is supplied (tests) it degrades to a single boot-time sweep so no ticker goroutine outlives the caller. The sweep stays age-guarded at MaxTokenTTL, so nothing reconstructable is ever reaped. Because the node sweep now runs during live traffic, it snapshots the live job set (Server.activeSessionIDs) and spares those dirs by id rather than by age alone — a long-lived session that only re-serves already-written segments stops advancing its dir mtime, which age could otherwise misclassify. In integrated mode the native and compat sweeps share one TranscodeDir but each snapshots only its own manager's live set; the resulting cross-manager reap of a >24h idle dir is bounded (rebuilds from token/recipe) and documented at both call sites. Part of #412
This commit is contained in:
@@ -665,6 +665,9 @@ func main() {
|
||||
// start, so this node can rebuild a Jellyfin transcode after its own
|
||||
// restart (the node hop token is recipe-less). Shares the offload Redis.
|
||||
srv.SetRecipeStore(noderecipe.NewStore(redisClient, 0))
|
||||
// Reclaim orphaned transcode dirs at boot and hourly thereafter, bound
|
||||
// to appCtx so it stops on shutdown.
|
||||
srv.StartOrphanSweeper(appCtx)
|
||||
handler = srv.Handler()
|
||||
}
|
||||
|
||||
@@ -2359,6 +2362,7 @@ func main() {
|
||||
if (mode == "integrated" || mode == "api") && cfg.JellyfinCompat.Enabled && cfg.JellyfinCompat.Listen != "" {
|
||||
compatDeps := jellycompat.Dependencies{
|
||||
Config: cfg,
|
||||
AppContext: appCtx,
|
||||
LiveConfig: configWatcher.Config,
|
||||
DB: deps.DB,
|
||||
SecretCipher: dataCipher,
|
||||
|
||||
@@ -857,11 +857,13 @@ func NewRouter(deps Dependencies) chi.Router {
|
||||
playbackHandler.PlaybackConfig = func() config.PlaybackConfig {
|
||||
return deps.CurrentConfig().Playback
|
||||
}
|
||||
if cleaned, err := playbackHandler.CleanupOrphanedTranscodes(); err != nil {
|
||||
slog.Warn("playback transcode cleanup failed", "dir", deps.Config.Playback.TranscodeDir, "error", err)
|
||||
} else if cleaned > 0 {
|
||||
slog.Info("playback transcode cleanup removed orphaned dirs", "dir", deps.Config.Playback.TranscodeDir, "count", cleaned)
|
||||
}
|
||||
// In integrated mode this and the jellycompat sweep both scan the same
|
||||
// TranscodeDir but each snapshots only its own manager's live set, so a
|
||||
// >24h idle dir owned by the other manager can be reaped. Bounded and
|
||||
// safe: active dirs stay mtime-fresh (spared) and either side rebuilds
|
||||
// from its token/recipe, so the worst case is a wasted rebuild. A shared
|
||||
// active-set source across both managers would remove even that.
|
||||
playback.StartPeriodicOrphanCleanup(deps.AppContext, "api", deps.Config.Playback.TranscodeDir, playbackHandler.CleanupOrphanedTranscodes, playback.OrphanCleanupInterval)
|
||||
}
|
||||
playbackHandler.ProbeEnsurer = deps.ProbeEnsurer
|
||||
playbackHandler.ChapterThumbnailQueuer = deps.ChapterThumbnailQueuer
|
||||
|
||||
@@ -111,11 +111,11 @@ func NewRouter(deps Dependencies) chi.Router {
|
||||
// Compat transcode reconstruct is driven by the recipe carried in the durable
|
||||
// compat playback store (jellycompat_playback_sessions); no separate native
|
||||
// recipe table is needed.
|
||||
if cleaned, err := playbackHandler.CleanupOrphanedTranscodes(); err != nil {
|
||||
slog.Warn("jellycompat transcode cleanup failed", "dir", playbackHandler.TranscodeDir, "error", err)
|
||||
} else if cleaned > 0 {
|
||||
slog.Info("jellycompat transcode cleanup removed orphaned dirs", "dir", playbackHandler.TranscodeDir, "count", cleaned)
|
||||
}
|
||||
//
|
||||
// This shares TranscodeDir with the native api sweep but snapshots only this
|
||||
// manager's live set; see the api NewRouter call site for why cross-manager
|
||||
// reaping of a >24h idle dir is bounded and safe.
|
||||
playback.StartPeriodicOrphanCleanup(deps.AppContext, "jellycompat", playbackHandler.TranscodeDir, playbackHandler.CleanupOrphanedTranscodes, playback.OrphanCleanupInterval)
|
||||
playbackHandler.profileRefreshRequester = deps.RecWorker
|
||||
playbackHandler.SettingsRepo = deps.SettingsRepo
|
||||
playbackHandler.RecipeNodeStore = deps.RecipeNodeStore
|
||||
|
||||
@@ -25,6 +25,10 @@ import (
|
||||
// Dependencies holds the pluggable pieces used by the compat server.
|
||||
type Dependencies struct {
|
||||
Config *config.Config
|
||||
// AppContext is the process lifecycle context. When set, it bounds the
|
||||
// periodic orphan-transcode sweep so it stops on shutdown; nil (tests) makes
|
||||
// the sweep a single boot-time run instead of a long-lived ticker.
|
||||
AppContext context.Context
|
||||
// LiveConfig returns the current hot-reloaded config. May be nil (tests,
|
||||
// worker modes); read through CurrentConfig(), which falls back to Config.
|
||||
LiveConfig func() *config.Config
|
||||
|
||||
@@ -1,13 +1,19 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"runtime/debug"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
var orphanCleanupMu sync.Mutex
|
||||
|
||||
// CleanupOrphanedTranscodeDirs removes per-session transcode directories that
|
||||
// are not associated with any currently active session IDs.
|
||||
//
|
||||
@@ -20,6 +26,11 @@ import (
|
||||
// Pass 0 to disable age-sparing (e.g. a dedicated node's boot-time full wipe,
|
||||
// where node restart is an accepted session loss).
|
||||
func CleanupOrphanedTranscodeDirs(root string, activeSessionIDs map[string]struct{}, minAge time.Duration) (int, error) {
|
||||
// Serialize concurrent orphan sweeps of the shared transcode root so two
|
||||
// rare startup sweeps cannot race on os.RemoveAll; one global mutex is fine.
|
||||
orphanCleanupMu.Lock()
|
||||
defer orphanCleanupMu.Unlock()
|
||||
|
||||
if root == "" {
|
||||
return 0, nil
|
||||
}
|
||||
@@ -70,6 +81,64 @@ func CleanupOrphanedTranscodeDirs(root string, activeSessionIDs map[string]struc
|
||||
return removed, nil
|
||||
}
|
||||
|
||||
// OrphanCleanupInterval is how often the periodic orphan sweep re-runs. A dir is
|
||||
// only reapable once it is older than MaxTokenTTL (24h), so sweeping much more
|
||||
// often than hourly buys nothing; hourly bounds how long an untracked orphan (a
|
||||
// dir whose owning session vanished without its RemoveAll succeeding) lingers on
|
||||
// a process that is never restarted, without adding meaningful load.
|
||||
const OrphanCleanupInterval = time.Hour
|
||||
|
||||
// runOrphanCleanup executes one sweep, recovering from a panic (a background
|
||||
// goroutine's unrecovered panic would crash the process) and logging the result.
|
||||
func runOrphanCleanup(component, dir string, cleanup func() (int, error)) {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
slog.Error("transcode cleanup panicked", "component", component, "dir", dir, "panic", r, "stack", string(debug.Stack()))
|
||||
}
|
||||
}()
|
||||
if cleaned, err := cleanup(); err != nil {
|
||||
slog.Warn("transcode cleanup failed", "component", component, "dir", dir, "error", err)
|
||||
} else if cleaned > 0 {
|
||||
slog.Info("transcode cleanup removed orphaned dirs", "component", component, "dir", dir, "count", cleaned)
|
||||
}
|
||||
}
|
||||
|
||||
// StartBackgroundOrphanCleanup runs a single orphaned-transcode sweep in its own
|
||||
// goroutine so a slow network-filesystem delete never blocks server startup.
|
||||
// The sweep is already safe to run concurrently with request handling: it
|
||||
// spares live/in-flight sessions and any dir younger than MaxTokenTTL, and
|
||||
// CleanupOrphanedTranscodeDirs serializes concurrent sweeps of the same root.
|
||||
func StartBackgroundOrphanCleanup(component, dir string, cleanup func() (int, error)) {
|
||||
go runOrphanCleanup(component, dir, cleanup)
|
||||
}
|
||||
|
||||
// StartPeriodicOrphanCleanup runs an immediate background sweep and then repeats
|
||||
// it every interval until ctx is cancelled. The startup sweep only reclaims dirs
|
||||
// orphaned by an ungraceful prior shutdown; the periodic re-run additionally
|
||||
// bounds "untracked orphan" accumulation (a dir whose owning session was dropped
|
||||
// without its RemoveAll succeeding) on a process that stays up for weeks. When
|
||||
// ctx is nil or interval is non-positive it degrades to a single boot-time sweep
|
||||
// so no ticker goroutine outlives a caller with no lifecycle handle (e.g. tests).
|
||||
func StartPeriodicOrphanCleanup(ctx context.Context, component, dir string, cleanup func() (int, error), interval time.Duration) {
|
||||
if ctx == nil || interval <= 0 {
|
||||
StartBackgroundOrphanCleanup(component, dir, cleanup)
|
||||
return
|
||||
}
|
||||
go func() {
|
||||
runOrphanCleanup(component, dir, cleanup)
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
runOrphanCleanup(component, dir, cleanup)
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func transcodeDirBelongsToActiveSession(name string, activeSessionIDs map[string]struct{}) bool {
|
||||
if _, ok := activeSessionIDs[name]; ok {
|
||||
return true
|
||||
|
||||
@@ -175,16 +175,52 @@ func NewServer(watcher *nodeconfig.Watcher, tracker *nodesessions.Tracker) *Serv
|
||||
sessions: make(map[string]*playback.TranscodeSession),
|
||||
lastAccess: make(map[string]time.Time),
|
||||
}
|
||||
if cfg := watcher.Config(); cfg != nil {
|
||||
if cleaned, err := playback.CleanupOrphanedTranscodeDirs(cfg.Playback.TranscodeDir, nil, 0); err != nil {
|
||||
slog.Warn("transcode node cleanup failed", "dir", cfg.Playback.TranscodeDir, "error", err)
|
||||
} else if cleaned > 0 {
|
||||
slog.Info("transcode node cleanup removed orphaned dirs", "dir", cfg.Playback.TranscodeDir, "count", cleaned)
|
||||
}
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
// StartOrphanSweeper runs the age-guarded orphan-transcode sweep immediately and
|
||||
// then hourly until ctx is cancelled. It never blocks (a slow network-filesystem
|
||||
// delete runs in its own goroutine), so it is safe to call before the node binds
|
||||
// its listener. This is the node's only filesystem-level reclaimer of dirs left
|
||||
// behind by a session that was dropped without its output dir being removed — the
|
||||
// idle reaper only deletes dirs it still tracks in s.sessions, so without this
|
||||
// periodic pass such orphans would linger until the next process restart. The
|
||||
// MaxTokenTTL age guard keeps a delete from racing a token-carried reconstruct
|
||||
// writing into TranscodeDir/<sessionID>: a dir younger than the max token
|
||||
// lifetime may still be reused, while older dirs are never reconstructable.
|
||||
func (s *Server) StartOrphanSweeper(ctx context.Context) {
|
||||
dir := ""
|
||||
if cfg := s.watcher.Config(); cfg != nil {
|
||||
dir = cfg.Playback.TranscodeDir
|
||||
}
|
||||
playback.StartPeriodicOrphanCleanup(ctx, "transcodenode", dir, func() (int, error) {
|
||||
// Re-read config each run so a hot-reloaded TranscodeDir is honored.
|
||||
cfg := s.watcher.Config()
|
||||
if cfg == nil {
|
||||
return 0, nil
|
||||
}
|
||||
// Spare the live registered jobs by id, not by age alone: now that the
|
||||
// sweep runs during live traffic, a long-lived session that re-serves
|
||||
// already-written segments stops advancing its dir mtime, so the age
|
||||
// guard could misclassify it as orphaned. The live set is authoritative
|
||||
// (in-flight reconstructs are covered by their fresh writes + age guard).
|
||||
return playback.CleanupOrphanedTranscodeDirs(cfg.Playback.TranscodeDir, s.activeSessionIDs(), playback.MaxTokenTTL)
|
||||
}, playback.OrphanCleanupInterval)
|
||||
}
|
||||
|
||||
// activeSessionIDs snapshots the ids of currently registered jobs so the orphan
|
||||
// sweep spares their output dirs regardless of directory mtime, mirroring the
|
||||
// central TranscodeManager's live-set snapshot.
|
||||
func (s *Server) activeSessionIDs() map[string]struct{} {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
active := make(map[string]struct{}, len(s.sessions))
|
||||
for id := range s.sessions {
|
||||
active[id] = struct{}{}
|
||||
}
|
||||
return active
|
||||
}
|
||||
|
||||
func (s *Server) SetFFmpegLogSink(sink playback.FFmpegLogSink) {
|
||||
s.ffmpegSink = sink
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user