From bfb0243712743e580159510ca7dd5379937af765 Mon Sep 17 00:00:00 2001 From: ShukeBta Date: Tue, 9 Jun 2026 19:18:21 +0800 Subject: [PATCH] fix: limit ffprobe concurrency --- internal/config/config.go | 25 ++++--- internal/handler/admin.go | 3 + internal/handler/system_extra.go | 1 + internal/service/ffprobe.go | 59 +++++++++++++++- internal/service/ffprobe_test.go | 100 ++++++++++++++++++++++----- internal/service/runtime_settings.go | 10 +++ web/src/pages/SettingsPage.tsx | 7 ++ 7 files changed, 174 insertions(+), 31 deletions(-) diff --git a/internal/config/config.go b/internal/config/config.go index 859c509..588bb6c 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -64,16 +64,20 @@ type TranscoderConfig struct { // AppConfig 保存运行时应用参数。 type AppConfig struct { - Port int `mapstructure:"port"` - Debug bool `mapstructure:"debug"` - Env string `mapstructure:"env"` - DataDir string `mapstructure:"data_dir"` - WebDir string `mapstructure:"web_dir"` - FFmpegPath string `mapstructure:"ffmpeg_path"` - FFprobePath string `mapstructure:"ffprobe_path"` - VAAPIDevice string `mapstructure:"vaapi_device"` - CORSOrigins []string `mapstructure:"cors_origins"` - ServerURL string `mapstructure:"server_url"` + Port int `mapstructure:"port"` + Debug bool `mapstructure:"debug"` + Env string `mapstructure:"env"` + DataDir string `mapstructure:"data_dir"` + WebDir string `mapstructure:"web_dir"` + FFmpegPath string `mapstructure:"ffmpeg_path"` + FFprobePath string `mapstructure:"ffprobe_path"` + // FFprobeMaxConcurrent limits concurrent ffprobe/ffmpeg metadata probes. + // NAS devices can become unresponsive when a scan starts many probe + // processes at once, so the default is deliberately conservative. + FFprobeMaxConcurrent int `mapstructure:"ffprobe_max_concurrent"` + VAAPIDevice string `mapstructure:"vaapi_device"` + CORSOrigins []string `mapstructure:"cors_origins"` + ServerURL string `mapstructure:"server_url"` } // DatabaseConfig 配置 GORM + SQLite。 @@ -213,6 +217,7 @@ func setDefaults(v *viper.Viper) { v.SetDefault("app.web_dir", "./web/dist") v.SetDefault("app.ffmpeg_path", "ffmpeg") v.SetDefault("app.ffprobe_path", "ffprobe") + v.SetDefault("app.ffprobe_max_concurrent", 1) v.SetDefault("app.vaapi_device", "/dev/dri/renderD128") v.SetDefault("app.cors_origins", []string{}) v.SetDefault("app.server_url", "") diff --git a/internal/handler/admin.go b/internal/handler/admin.go index 290fe1c..f516b8d 100644 --- a/internal/handler/admin.go +++ b/internal/handler/admin.go @@ -260,6 +260,9 @@ func updateSettingHandler(svc *service.Container) gin.HandlerFunc { _ = svc.Repo.DB.WithContext(c.Request.Context()).Model(&model.User{}).Where("hide_adult = ?", false).Update("hide_adult", true).Error } service.ApplyRuntimeSetting(svc.Cfg, req.Key, req.Value) + if svc.FFprobe != nil && (req.Key == "ffprobe.max_concurrent" || req.Key == "app.ffprobe_max_concurrent") { + svc.FFprobe.SetMaxConcurrent(svc.Cfg.App.FFprobeMaxConcurrent) + } if req.Key == "transcode.enabled" && !svc.Cfg.Transcoder.Enabled { svc.Transcoder.StopAll() } diff --git a/internal/handler/system_extra.go b/internal/handler/system_extra.go index 17bf2e7..1fe5794 100644 --- a/internal/handler/system_extra.go +++ b/internal/handler/system_extra.go @@ -72,6 +72,7 @@ func schemaHandler(_ *service.Container) gin.HandlerFunc { {"key": "transcode.idle_timeout_seconds", "type": "number", "label": "转码空闲停止秒数"}, {"key": "ffmpeg.path", "type": "text", "label": "FFmpeg 路径"}, {"key": "ffprobe.path", "type": "text", "label": "FFprobe 路径"}, + {"key": "ffprobe.max_concurrent", "type": "number", "label": "FFprobe 最大并发"}, }, }, { diff --git a/internal/service/ffprobe.go b/internal/service/ffprobe.go index 3922229..bb75287 100644 --- a/internal/service/ffprobe.go +++ b/internal/service/ffprobe.go @@ -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 { diff --git a/internal/service/ffprobe_test.go b/internal/service/ffprobe_test.go index 7404f30..160f952 100644 --- a/internal/service/ffprobe_test.go +++ b/internal/service/ffprobe_test.go @@ -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) } } diff --git a/internal/service/runtime_settings.go b/internal/service/runtime_settings.go index fb77cca..9818f6b 100644 --- a/internal/service/runtime_settings.go +++ b/internal/service/runtime_settings.go @@ -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": diff --git a/web/src/pages/SettingsPage.tsx b/web/src/pages/SettingsPage.tsx index 5ca9d4c..1f4f14a 100644 --- a/web/src/pages/SettingsPage.tsx +++ b/web/src/pages/SettingsPage.tsx @@ -121,6 +121,13 @@ const GROUPS: SettingGroup[] = [ type: 'text', placeholder: 'ffprobe', }, + { + key: 'ffprobe.max_concurrent', + label: 'FFprobe 最大并发', + type: 'number', + hint: 'NAS 建议 1;用于扫描、整理洗版和手动探测,避免同时启动多个 ffprobe 进程', + defaultValue: '1', + }, ], }, {