package service import ( "context" "os" "time" "go.uber.org/zap" ) // WaitReady blocks (with a deadline) until the playlist file shows up on // disk for the *current* job generation. Stale playlists left behind by a // failed RemoveAll / still-exiting ffmpeg must not unblock a mid-file restart. func (t *TranscoderService) WaitReady(ctx context.Context, mediaID string, timeout time.Duration) bool { deadline := time.Now().Add(timeout) for { t.mu.Lock() job, ok := t.jobs[mediaID] var started time.Time if ok { started = job.startedAt } t.mu.Unlock() if ok { if info, err := os.Stat(t.PlaylistPath(mediaID)); err == nil { // Allow a small clock skew; reject anything older than this job. if !info.ModTime().Before(started.Add(-2 * time.Second)) { t.mu.Lock() if j, exists := t.jobs[mediaID]; exists { j.playlistOK = true } t.mu.Unlock() return true } } } if time.Now().After(deadline) || ctx.Err() != nil { return false } select { case <-ctx.Done(): return false case <-time.After(250 * time.Millisecond): } } } // StopJob cancels a running ffmpeg process for mediaID, if any, and waits // briefly for it to exit so a subsequent EnsureJobFrom cannot race-write the // same HLS directory. func (t *TranscoderService) StopJob(mediaID string) { gate := t.mediaStartGate(mediaID) gate.Lock() defer gate.Unlock() t.mu.Lock() prev := t.detachJobLocked(mediaID) t.mu.Unlock() waitJobExit(prev, 12*time.Second) } // TouchJob records client activity for the HLS playlist or segment. The idle // watchdog uses it to stop ffmpeg soon after the player is closed or switches // back to direct play. func (t *TranscoderService) TouchJob(mediaID string) { t.mu.Lock() defer t.mu.Unlock() t.touchJobLocked(mediaID) } func (t *TranscoderService) touchJobLocked(mediaID string) { if j, ok := t.jobs[mediaID]; ok { j.lastAccess = time.Now() } } // StopAll terminates every running transcode (called on graceful shutdown). func (t *TranscoderService) StopAll() { t.mu.Lock() pending := make([]*hlsJob, 0, len(t.jobs)) for id := range t.jobs { pending = append(pending, t.detachJobLocked(id)) } t.mu.Unlock() for _, j := range pending { waitJobExit(j, 5*time.Second) } } // ActiveJob is the JSON shape exposed to the React Tasks panel. type ActiveJob struct { MediaID string `json:"media_id"` Encoder string `json:"encoder"` StartedAt time.Time `json:"started_at"` PlaylistOK bool `json:"playlist_ok"` } // Active returns a snapshot of the currently running transcode jobs. func (t *TranscoderService) Active() []ActiveJob { t.mu.Lock() defer t.mu.Unlock() out := make([]ActiveJob, 0, len(t.jobs)) for _, j := range t.jobs { out = append(out, ActiveJob{ MediaID: j.mediaID, Encoder: j.encoder, StartedAt: j.startedAt, PlaylistOK: j.playlistOK, }) } return out } func (t *TranscoderService) maxConcurrent() int { if t.cfg.Transcoder.MaxConcurrent <= 0 { return 1 } return t.cfg.Transcoder.MaxConcurrent } func (t *TranscoderService) idleTimeout() time.Duration { if t.cfg.Transcoder.IdleTimeoutSeconds <= 0 { return 120 * time.Second } return time.Duration(t.cfg.Transcoder.IdleTimeoutSeconds) * time.Second } func (t *TranscoderService) monitorIdle(ctx context.Context, job *hlsJob) { timeout := t.idleTimeout() ticker := time.NewTicker(15 * time.Second) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: t.mu.Lock() current, ok := t.jobs[job.mediaID] if !ok { t.mu.Unlock() return } idleFor := time.Since(current.lastAccess) t.mu.Unlock() if idleFor >= timeout { t.log.Info("transcode idle timeout", zap.String("media_id", job.mediaID), zap.Duration("idle_for", idleFor), zap.Duration("timeout", timeout), ) t.StopJob(job.mediaID) return } } } }