mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-07 13:56:37 +08:00
split transcoder service helpers
This commit is contained in:
@@ -25,11 +25,8 @@ package service
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
|
||||||
"os"
|
"os"
|
||||||
"os/exec"
|
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -144,269 +141,3 @@ func (t *TranscoderService) EnsureJob(ctx context.Context, mediaID string) (stri
|
|||||||
go t.runFFmpeg(jobCtx, job, m.Path)
|
go t.runFFmpeg(jobCtx, job, m.Path)
|
||||||
return t.PlaylistPath(mediaID), nil
|
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"))
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -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"))
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user