diff --git a/cmd/silo/main.go b/cmd/silo/main.go index b0939170..38c130ba 100644 --- a/cmd/silo/main.go +++ b/cmd/silo/main.go @@ -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, diff --git a/internal/api/router.go b/internal/api/router.go index 61262c44..535a3d53 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -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 diff --git a/internal/jellycompat/router.go b/internal/jellycompat/router.go index 8f3f0d54..31a3b0c2 100644 --- a/internal/jellycompat/router.go +++ b/internal/jellycompat/router.go @@ -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 diff --git a/internal/jellycompat/server.go b/internal/jellycompat/server.go index 0e70a0ca..b1c5fab2 100644 --- a/internal/jellycompat/server.go +++ b/internal/jellycompat/server.go @@ -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 diff --git a/internal/playback/transcode_cleanup.go b/internal/playback/transcode_cleanup.go index 259d24e9..1cae53a5 100644 --- a/internal/playback/transcode_cleanup.go +++ b/internal/playback/transcode_cleanup.go @@ -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 diff --git a/internal/transcodenode/server.go b/internal/transcodenode/server.go index 02cdbfc9..0da2af94 100644 --- a/internal/transcodenode/server.go +++ b/internal/transcodenode/server.go @@ -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/: 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 }