fix: limit ffprobe concurrency

This commit is contained in:
ShukeBta
2026-06-09 19:18:21 +08:00
parent 7d7f3cc758
commit bfb0243712
7 changed files with 174 additions and 31 deletions
+56 -3
View File
@@ -16,6 +16,7 @@ import (
"regexp"
"strconv"
"strings"
"sync"
"time"
"go.uber.org/zap"
@@ -25,13 +26,35 @@ import (
// FFprobeService wraps the external ffprobe binary.
type FFprobeService struct {
cfg *config.Config
log *zap.Logger
cfg *config.Config
log *zap.Logger
mu sync.RWMutex
limiter chan struct{}
}
// NewFFprobeService is the constructor.
func NewFFprobeService(cfg *config.Config, log *zap.Logger) *FFprobeService {
return &FFprobeService{cfg: cfg, log: log}
maxConcurrent := normalizeFFprobeMaxConcurrent(cfg.App.FFprobeMaxConcurrent)
return &FFprobeService{cfg: cfg, log: log, limiter: make(chan struct{}, maxConcurrent)}
}
func normalizeFFprobeMaxConcurrent(n int) int {
if n <= 0 {
return 1
}
if n > 8 {
return 8
}
return n
}
func (f *FFprobeService) SetMaxConcurrent(n int) {
if f == nil {
return
}
f.mu.Lock()
defer f.mu.Unlock()
f.limiter = make(chan struct{}, normalizeFFprobeMaxConcurrent(n))
}
// ProbeResult is the subset of ffprobe output consumed by the scanner.
@@ -50,6 +73,11 @@ func (f *FFprobeService) Probe(ctx context.Context, path string) (*ProbeResult,
if f == nil {
return nil, errors.New("ffprobe service nil")
}
token, err := f.acquire(ctx)
if err != nil {
return nil, err
}
defer f.release(token)
if bin, err := resolveLocalExecutable(f.cfg.App.FFprobePath, "ffprobe"); err == nil {
f.cfg.App.FFprobePath = bin
probeCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
@@ -73,6 +101,31 @@ func (f *FFprobeService) Probe(ctx context.Context, path string) (*ProbeResult,
return f.probeWithFFmpeg(ctx, path)
}
func (f *FFprobeService) acquire(ctx context.Context) (chan struct{}, error) {
f.mu.RLock()
limiter := f.limiter
f.mu.RUnlock()
if limiter == nil {
return nil, nil
}
select {
case limiter <- struct{}{}:
return limiter, nil
case <-ctx.Done():
return nil, ctx.Err()
}
}
func (f *FFprobeService) release(limiter chan struct{}) {
if limiter == nil {
return
}
select {
case <-limiter:
default:
}
}
func (f *FFprobeService) probeWithFFmpeg(ctx context.Context, path string) (*ProbeResult, error) {
bin, err := resolveLocalExecutable(f.cfg.App.FFmpegPath, "ffmpeg")
if err != nil {
+82 -18
View File
@@ -1,24 +1,88 @@
package service
import "testing"
import (
"context"
"testing"
"time"
func TestParseFFmpegProbeText(t *testing.T) {
text := `Input #0, matroska,webm, from 'show.mkv':
Duration: 00:23:42.11, start: 0.000000, bitrate: 5132 kb/s
Stream #0:0: Video: h264 (Main), yuv420p(progressive), 1920x1080 [SAR 1:1 DAR 16:9], 23.98 fps
Stream #0:1(jpn): Audio: eac3, 48000 Hz, stereo, fltp, 128 kb/s (default)`
"go.uber.org/zap"
got := parseFFmpegProbeText(text)
if got.Container != "matroska,webm" {
t.Fatalf("container = %q", got.Container)
}
if got.DurationSec != 1422 {
t.Fatalf("duration = %d", got.DurationSec)
}
if got.VideoCodec != "h264" || got.Width != 1920 || got.Height != 1080 {
t.Fatalf("video = %#v", got)
}
if got.AudioCodec != "eac3" {
t.Fatalf("audio = %q", got.AudioCodec)
"github.com/ShukeBta/MediaStationGo/internal/config"
)
func TestFFprobeServiceDefaultsToSingleConcurrentProbe(t *testing.T) {
svc := NewFFprobeService(&config.Config{}, zap.NewNop())
if got := cap(svc.limiter); got != 1 {
t.Fatalf("limiter capacity = %d, want 1", got)
}
}
func TestFFprobeServiceClampsConfiguredConcurrency(t *testing.T) {
cfg := &config.Config{}
cfg.App.FFprobeMaxConcurrent = 3
svc := NewFFprobeService(cfg, zap.NewNop())
if got := cap(svc.limiter); got != 3 {
t.Fatalf("limiter capacity = %d, want 3", got)
}
cfg.App.FFprobeMaxConcurrent = 99
svc = NewFFprobeService(cfg, zap.NewNop())
if got := cap(svc.limiter); got != 8 {
t.Fatalf("limiter capacity = %d, want max clamp 8", got)
}
}
func TestFFprobeAcquireHonorsContextWhenLimitReached(t *testing.T) {
svc := &FFprobeService{limiter: make(chan struct{}, 1)}
token, err := svc.acquire(t.Context())
if err != nil {
t.Fatalf("first acquire: %v", err)
}
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Millisecond)
defer cancel()
if _, err := svc.acquire(ctx); err == nil {
t.Fatal("second acquire should block until context deadline")
}
svc.release(token)
token, err = svc.acquire(t.Context())
if err != nil {
t.Fatalf("acquire after release: %v", err)
}
svc.release(token)
}
func TestFFprobeSetMaxConcurrentHotSwapsLimiter(t *testing.T) {
svc := NewFFprobeService(&config.Config{}, zap.NewNop())
firstToken, err := svc.acquire(t.Context())
if err != nil {
t.Fatalf("first acquire: %v", err)
}
svc.SetMaxConcurrent(2)
secondToken, err := svc.acquire(t.Context())
if err != nil {
t.Fatalf("acquire after resize: %v", err)
}
thirdToken, err := svc.acquire(t.Context())
if err != nil {
t.Fatalf("second acquire after resize: %v", err)
}
svc.release(firstToken)
svc.release(secondToken)
svc.release(thirdToken)
}
func TestApplyRuntimeSettingFFprobeMaxConcurrent(t *testing.T) {
cfg := &config.Config{}
ApplyRuntimeSetting(cfg, "ffprobe.max_concurrent", "4")
if cfg.App.FFprobeMaxConcurrent != 4 {
t.Fatalf("FFprobeMaxConcurrent = %d, want 4", cfg.App.FFprobeMaxConcurrent)
}
ApplyRuntimeSetting(cfg, "ffprobe.max_concurrent", "99")
if cfg.App.FFprobeMaxConcurrent != 8 {
t.Fatalf("FFprobeMaxConcurrent = %d, want clamp 8", cfg.App.FFprobeMaxConcurrent)
}
ApplyRuntimeSetting(cfg, "ffprobe.max_concurrent", "0")
if cfg.App.FFprobeMaxConcurrent != 1 {
t.Fatalf("FFprobeMaxConcurrent = %d, want clamp 1", cfg.App.FFprobeMaxConcurrent)
}
}
+10
View File
@@ -37,6 +37,16 @@ func ApplyRuntimeSetting(cfg *config.Config, key, value string) {
cfg.App.FFmpegPath = value
case "ffprobe.path", "app.ffprobe_path":
cfg.App.FFprobePath = value
case "ffprobe.max_concurrent", "app.ffprobe_max_concurrent":
if n, err := strconv.Atoi(value); err == nil {
if n < 1 {
n = 1
}
if n > 8 {
n = 8
}
cfg.App.FFprobeMaxConcurrent = n
}
case "transcode.enabled", "transcoder.enabled":
cfg.Transcoder.Enabled = parseBoolSetting(value, true)
case "transcode.hw_enabled", "transcoder.hardware_accel":