From 008f3f233632bf607b50a5eb25958784a5f23b56 Mon Sep 17 00:00:00 2001 From: ShukeBta <272197458+ShukeBta@users.noreply.github.com> Date: Sat, 27 Jun 2026 02:42:42 +0800 Subject: [PATCH] split transcoder service helpers --- internal/service/transcoder.go | 269 -------------------- internal/service/transcoder_ffmpeg_probe.go | 152 +++++++++++ internal/service/transcoder_jobs.go | 137 ++++++++++ 3 files changed, 289 insertions(+), 269 deletions(-) create mode 100644 internal/service/transcoder_ffmpeg_probe.go create mode 100644 internal/service/transcoder_jobs.go diff --git a/internal/service/transcoder.go b/internal/service/transcoder.go index 675d13b..1f29974 100644 --- a/internal/service/transcoder.go +++ b/internal/service/transcoder.go @@ -25,11 +25,8 @@ package service import ( "context" "errors" - "fmt" "os" - "os/exec" "path/filepath" - "strings" "sync" "time" @@ -144,269 +141,3 @@ func (t *TranscoderService) EnsureJob(ctx context.Context, mediaID string) (stri go t.runFFmpeg(jobCtx, job, m.Path) return t.PlaylistPath(mediaID), nil } - -// WaitReady blocks (with a deadline) until the playlist file shows up on -// disk. Returns true on success. -func (t *TranscoderService) WaitReady(ctx context.Context, mediaID string, timeout time.Duration) bool { - deadline := time.Now().Add(timeout) - for { - if _, err := os.Stat(t.PlaylistPath(mediaID)); err == nil { - t.mu.Lock() - if j, ok := t.jobs[mediaID]; ok { - 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. -func (t *TranscoderService) StopJob(mediaID string) { - t.mu.Lock() - defer t.mu.Unlock() - if j, ok := t.jobs[mediaID]; ok { - j.cancel() - delete(t.jobs, mediaID) - } -} - -// 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() - defer t.mu.Unlock() - for id, j := range t.jobs { - j.cancel() - delete(t.jobs, id) - } -} - -// 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 - } - } - } -} - -func (t *TranscoderService) runFFmpeg(ctx context.Context, job *hlsJob, source string) { - bin, err := t.resolveFFmpegPath() - if err != nil { - t.log.Warn("ffmpeg unavailable", zap.String("media_id", job.mediaID), zap.Error(err)) - t.mu.Lock() - delete(t.jobs, job.mediaID) - t.mu.Unlock() - t.hub.Publish("transcode", map[string]any{ - "media_id": job.mediaID, - "status": "error", - "error": err.Error(), - }) - return - } - - playlist := filepath.Join(job.outputDir, "index.m3u8") - segments := filepath.Join(job.outputDir, "seg_%05d.ts") - - args := buildFFmpegArgs(t.cfg, source, playlist, segments) - - cmd := exec.CommandContext(ctx, bin, args...) // #nosec G204 -- bin is resolved by resolveFFmpegPath and args are passed without a shell. - cmd.Stderr = os.Stderr - - t.log.Info("transcode started", - zap.String("media_id", job.mediaID), - zap.String("encoder", job.encoder), - zap.String("source", source), - ) - t.hub.Publish("transcode", map[string]any{ - "media_id": job.mediaID, - "encoder": job.encoder, - "status": "started", - }) - - if err := cmd.Run(); err != nil && !errors.Is(ctx.Err(), context.Canceled) { - t.log.Warn("ffmpeg exited", - zap.String("media_id", job.mediaID), - zap.Error(err), - ) - } - - t.mu.Lock() - delete(t.jobs, job.mediaID) - t.mu.Unlock() - - t.hub.Publish("transcode", map[string]any{ - "media_id": job.mediaID, - "status": "stopped", - "duration": time.Since(job.startedAt).Seconds(), - }) -} - -func (t *TranscoderService) resolveFFmpegPath() (string, error) { - var lastErr error - for _, bin := range executableCandidates(strings.TrimSpace(t.cfg.App.FFmpegPath), "ffmpeg") { - if err := validateFFmpegForTranscode(context.Background(), bin, t.effectiveEncoder()); err != nil { - lastErr = err - continue - } - t.cfg.App.FFmpegPath = bin - return bin, nil - } - if lastErr != nil { - return "", fmt.Errorf("no usable ffmpeg found for HLS transcode: %w", lastErr) - } - return "", fmt.Errorf("ffmpeg not found in PATH or common local app directories; configure app.ffmpeg_path to an existing local ffmpeg") -} - -func (t *TranscoderService) effectiveEncoder() string { - if !t.cfg.Transcoder.HardwareAccel { - return "" - } - return normalizedHardwareEncoder(t.cfg.Transcoder.Encoder) -} - -func normalizedHardwareEncoder(encoder string) string { - switch strings.ToLower(strings.TrimSpace(encoder)) { - case "nvenc", "qsv", "vaapi": - return strings.ToLower(strings.TrimSpace(encoder)) - default: - return "" - } -} - -func validateFFmpegForTranscode(ctx context.Context, bin, encoder string) error { - required := requiredVideoEncoder(encoder) - out, err := commandOutput(ctx, 8*time.Second, bin, "-hide_banner", "-encoders") - if err != nil { - return fmt.Errorf("%s cannot list encoders: %w", bin, err) - } - if !hasFFmpegListEntry(string(out), required) { - return fmt.Errorf("%s does not provide required encoder %s", bin, required) - } - - out, err = commandOutput(ctx, 8*time.Second, bin, "-hide_banner", "-muxers") - if err != nil || !hasFFmpegListEntry(string(out), "hls") { - out, err = commandOutput(ctx, 8*time.Second, bin, "-hide_banner", "-formats") - if err != nil { - return fmt.Errorf("%s cannot list muxers/formats: %w", bin, err) - } - } - if !hasFFmpegListEntry(string(out), "hls") { - return fmt.Errorf("%s does not provide hls muxer", bin) - } - return nil -} - -func requiredVideoEncoder(encoder string) string { - switch encoder { - case "nvenc": - return "h264_nvenc" - case "qsv": - return "h264_qsv" - case "vaapi": - return "h264_vaapi" - default: - return "libx264" - } -} - -func hasFFmpegListEntry(output, name string) bool { - for _, line := range strings.Split(output, "\n") { - fields := strings.Fields(line) - for _, field := range fields { - if field == name { - return true - } - } - } - return false -} - -// HumanFFmpegProfile is exposed for the admin UI / settings view. -func (t *TranscoderService) HumanFFmpegProfile() string { - return fmt.Sprintf("ffmpeg=%s, encoder=%s, output=%s", - t.cfg.App.FFmpegPath, t.cfg.Transcoder.Encoder, filepath.Join(t.cfg.Cache.CacheDir, "hls")) -} diff --git a/internal/service/transcoder_ffmpeg_probe.go b/internal/service/transcoder_ffmpeg_probe.go new file mode 100644 index 0000000..5d4800c --- /dev/null +++ b/internal/service/transcoder_ffmpeg_probe.go @@ -0,0 +1,152 @@ +package service + +import ( + "context" + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "strings" + "time" + + "go.uber.org/zap" +) + +func (t *TranscoderService) runFFmpeg(ctx context.Context, job *hlsJob, source string) { + bin, err := t.resolveFFmpegPath() + if err != nil { + t.log.Warn("ffmpeg unavailable", zap.String("media_id", job.mediaID), zap.Error(err)) + t.mu.Lock() + delete(t.jobs, job.mediaID) + t.mu.Unlock() + t.hub.Publish("transcode", map[string]any{ + "media_id": job.mediaID, + "status": "error", + "error": err.Error(), + }) + return + } + + playlist := filepath.Join(job.outputDir, "index.m3u8") + segments := filepath.Join(job.outputDir, "seg_%05d.ts") + + args := buildFFmpegArgs(t.cfg, source, playlist, segments) + + cmd := exec.CommandContext(ctx, bin, args...) // #nosec G204 -- bin is resolved by resolveFFmpegPath and args are passed without a shell. + cmd.Stderr = os.Stderr + + t.log.Info("transcode started", + zap.String("media_id", job.mediaID), + zap.String("encoder", job.encoder), + zap.String("source", source), + ) + t.hub.Publish("transcode", map[string]any{ + "media_id": job.mediaID, + "encoder": job.encoder, + "status": "started", + }) + + if err := cmd.Run(); err != nil && !errors.Is(ctx.Err(), context.Canceled) { + t.log.Warn("ffmpeg exited", + zap.String("media_id", job.mediaID), + zap.Error(err), + ) + } + + t.mu.Lock() + delete(t.jobs, job.mediaID) + t.mu.Unlock() + + t.hub.Publish("transcode", map[string]any{ + "media_id": job.mediaID, + "status": "stopped", + "duration": time.Since(job.startedAt).Seconds(), + }) +} + +func (t *TranscoderService) resolveFFmpegPath() (string, error) { + var lastErr error + for _, bin := range executableCandidates(strings.TrimSpace(t.cfg.App.FFmpegPath), "ffmpeg") { + if err := validateFFmpegForTranscode(context.Background(), bin, t.effectiveEncoder()); err != nil { + lastErr = err + continue + } + t.cfg.App.FFmpegPath = bin + return bin, nil + } + if lastErr != nil { + return "", fmt.Errorf("no usable ffmpeg found for HLS transcode: %w", lastErr) + } + return "", fmt.Errorf("ffmpeg not found in PATH or common local app directories; configure app.ffmpeg_path to an existing local ffmpeg") +} + +func (t *TranscoderService) effectiveEncoder() string { + if !t.cfg.Transcoder.HardwareAccel { + return "" + } + return normalizedHardwareEncoder(t.cfg.Transcoder.Encoder) +} + +func normalizedHardwareEncoder(encoder string) string { + switch strings.ToLower(strings.TrimSpace(encoder)) { + case "nvenc", "qsv", "vaapi": + return strings.ToLower(strings.TrimSpace(encoder)) + default: + return "" + } +} + +func validateFFmpegForTranscode(ctx context.Context, bin, encoder string) error { + required := requiredVideoEncoder(encoder) + out, err := commandOutput(ctx, 8*time.Second, bin, "-hide_banner", "-encoders") + if err != nil { + return fmt.Errorf("%s cannot list encoders: %w", bin, err) + } + if !hasFFmpegListEntry(string(out), required) { + return fmt.Errorf("%s does not provide required encoder %s", bin, required) + } + + out, err = commandOutput(ctx, 8*time.Second, bin, "-hide_banner", "-muxers") + if err != nil || !hasFFmpegListEntry(string(out), "hls") { + out, err = commandOutput(ctx, 8*time.Second, bin, "-hide_banner", "-formats") + if err != nil { + return fmt.Errorf("%s cannot list muxers/formats: %w", bin, err) + } + } + if !hasFFmpegListEntry(string(out), "hls") { + return fmt.Errorf("%s does not provide hls muxer", bin) + } + return nil +} + +func requiredVideoEncoder(encoder string) string { + switch encoder { + case "nvenc": + return "h264_nvenc" + case "qsv": + return "h264_qsv" + case "vaapi": + return "h264_vaapi" + default: + return "libx264" + } +} + +func hasFFmpegListEntry(output, name string) bool { + for _, line := range strings.Split(output, "\n") { + fields := strings.Fields(line) + for _, field := range fields { + if field == name { + return true + } + } + } + return false +} + +// HumanFFmpegProfile is exposed for the admin UI / settings view. +func (t *TranscoderService) HumanFFmpegProfile() string { + return fmt.Sprintf("ffmpeg=%s, encoder=%s, output=%s", + t.cfg.App.FFmpegPath, t.cfg.Transcoder.Encoder, filepath.Join(t.cfg.Cache.CacheDir, "hls")) +} diff --git a/internal/service/transcoder_jobs.go b/internal/service/transcoder_jobs.go new file mode 100644 index 0000000..5902e5d --- /dev/null +++ b/internal/service/transcoder_jobs.go @@ -0,0 +1,137 @@ +package service + +import ( + "context" + "os" + "time" + + "go.uber.org/zap" +) + +// WaitReady blocks (with a deadline) until the playlist file shows up on +// disk. Returns true on success. +func (t *TranscoderService) WaitReady(ctx context.Context, mediaID string, timeout time.Duration) bool { + deadline := time.Now().Add(timeout) + for { + if _, err := os.Stat(t.PlaylistPath(mediaID)); err == nil { + t.mu.Lock() + if j, ok := t.jobs[mediaID]; ok { + 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. +func (t *TranscoderService) StopJob(mediaID string) { + t.mu.Lock() + defer t.mu.Unlock() + if j, ok := t.jobs[mediaID]; ok { + j.cancel() + delete(t.jobs, mediaID) + } +} + +// 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() + defer t.mu.Unlock() + for id, j := range t.jobs { + j.cancel() + delete(t.jobs, id) + } +} + +// 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 + } + } + } +}